diff --git a/dot/network/interfaces.go b/dot/network/interfaces.go index acff76d6ea..6f9145367f 100644 --- a/dot/network/interfaces.go +++ b/dot/network/interfaces.go @@ -18,3 +18,9 @@ type Logger interface { Warnf(format string, args ...interface{}) Errorf(format string, args ...interface{}) } + +// MDNS is the mDNS service interface. +type MDNS interface { + Start() error + Stop() error +} diff --git a/dot/network/mdns.go b/dot/network/mdns.go deleted file mode 100644 index 97a08e2575..0000000000 --- a/dot/network/mdns.go +++ /dev/null @@ -1,95 +0,0 @@ -// Copyright 2021 ChainSafe Systems (ON) -// SPDX-License-Identifier: LGPL-3.0-only - -package network - -import ( - "context" - "time" - - "github.com/ChainSafe/gossamer/internal/log" - "github.com/libp2p/go-libp2p-core/peer" - "github.com/libp2p/go-libp2p-core/peerstore" - libp2pdiscovery "github.com/libp2p/go-libp2p/p2p/discovery/mdns_legacy" -) - -// MDNSPeriod is 1 minute -const MDNSPeriod = time.Minute - -// Notifee See https://godoc.org/github.com/libp2p/go-libp2p/p2p/discovery#Notifee -type Notifee struct { - logger Logger - ctx context.Context - host *host -} - -// mdns submodule -type mdns struct { - logger Logger - host *host - mdns libp2pdiscovery.Service -} - -// newMDNS creates a new mDNS instance from the host -func newMDNS(host *host) *mdns { - return &mdns{ - logger: log.NewFromGlobal(log.AddContext("module", "mdns")), - host: host, - } -} - -// startMDNS starts a new mDNS discovery service -func (m *mdns) start() { - m.logger.Debugf( - "Starting mDNS discovery service with host %s, period %s and protocol %s...", - m.host.id(), MDNSPeriod, m.host.protocolID) - - // create and start service - mdns, err := libp2pdiscovery.NewMdnsService( - m.host.ctx, - m.host.p2pHost, - MDNSPeriod, - string(m.host.protocolID), - ) - if err != nil { - m.logger.Errorf("Failed to start mDNS discovery service: %s", err) - return - } - - // register Notifee on service - mdns.RegisterNotifee(Notifee{ - logger: m.logger, - ctx: m.host.ctx, - host: m.host, - }) - - m.mdns = mdns -} - -// close shuts down the mDNS discovery service -func (m *mdns) close() error { - // check if service is running - if m.mdns == nil { - return nil - } - - // close service - err := m.mdns.Close() - if err != nil { - m.logger.Warnf("Failed to close mDNS discovery service: %s", err) - return err - } - - return nil -} - -// HandlePeerFound is event handler called when a peer is found -func (n Notifee) HandlePeerFound(p peer.AddrInfo) { - n.logger.Debugf( - "Peer %s found using mDNS discovery, with host %s", - p.ID, n.host.id()) - - n.host.p2pHost.Peerstore().AddAddrs(p.ID, p.Addrs, peerstore.PermanentAddrTTL) - // connect to found peer - n.host.cm.peerSetHandler.AddPeer(0, p.ID) -} diff --git a/dot/network/service.go b/dot/network/service.go index 93df156132..5797453458 100644 --- a/dot/network/service.go +++ b/dot/network/service.go @@ -6,6 +6,7 @@ package network import ( "context" "errors" + "fmt" "math/big" "strings" "sync" @@ -14,6 +15,7 @@ import ( "github.com/ChainSafe/gossamer/dot/peerset" "github.com/ChainSafe/gossamer/dot/telemetry" "github.com/ChainSafe/gossamer/internal/log" + "github.com/ChainSafe/gossamer/internal/mdns" "github.com/ChainSafe/gossamer/internal/metrics" "github.com/ChainSafe/gossamer/lib/common" libp2pnetwork "github.com/libp2p/go-libp2p-core/network" @@ -103,7 +105,7 @@ type Service struct { cfg *Config host *host - mdns *mdns + mdns MDNS gossip *gossip bufPool *sync.Pool streamManager *streamManager @@ -186,12 +188,20 @@ func NewService(cfg *Config) (*Service, error) { }, } + serviceTag := string(host.protocolID) + notifee := mdns.NewNotifeeTracker(host.p2pHost.Peerstore(), host.cm.peerSetHandler) + mdnsLogger := log.NewFromGlobal(log.AddContext("module", "mdns")) + mdnsLogger.Debugf( + "Creating mDNS discovery service with host %s and protocol %s...", + host.id(), host.protocolID) + mdnsService := mdns.NewService(host.p2pHost, serviceTag, mdnsLogger, notifee) + network := &Service{ ctx: ctx, cancel: cancel, cfg: cfg, host: host, - mdns: newMDNS(host), + mdns: mdnsService, gossip: newGossip(), blockState: cfg.BlockState, transactionHandler: cfg.TransactionHandler, @@ -303,7 +313,10 @@ func (s *Service) Start() error { s.startPeerSetHandler() if !s.noMDNS { - s.mdns.start() + err = s.mdns.Start() + if err != nil { + return fmt.Errorf("starting mDNS service: %w", err) + } } if !s.noDiscover { @@ -443,7 +456,7 @@ func (s *Service) Stop() error { s.cancel() // close mDNS discovery service - err := s.mdns.close() + err := s.mdns.Stop() if err != nil { logger.Errorf("Failed to close mDNS discovery service: %s", err) } diff --git a/go.mod b/go.mod index b41672efa7..1746d8d60b 100644 --- a/go.mod +++ b/go.mod @@ -119,7 +119,6 @@ require ( github.com/libp2p/go-openssl v0.0.7 // indirect github.com/libp2p/go-reuseport v0.2.0 // indirect github.com/libp2p/go-yamux/v3 v3.1.2 // indirect - github.com/libp2p/zeroconf/v2 v2.1.1 // indirect github.com/lucas-clemente/quic-go v0.27.1 // indirect github.com/marten-seemann/qtls-go1-16 v0.1.5 // indirect github.com/marten-seemann/qtls-go1-17 v0.1.1 // indirect @@ -170,7 +169,7 @@ require ( github.com/tomasen/realip v0.0.0-20180522021738-f0c99a92ddce // indirect github.com/vedhavyas/go-subkey v1.0.3 // indirect github.com/whyrusleeping/go-keyspace v0.0.0-20160322163242-5b898ac5add1 // indirect - github.com/whyrusleeping/mdns v0.0.0-20190826153040-b9b60ed33aa9 // indirect + github.com/whyrusleeping/mdns v0.0.0-20190826153040-b9b60ed33aa9 github.com/whyrusleeping/multiaddr-filter v0.0.0-20160516205228-e903e4adabd7 // indirect go.opencensus.io v0.23.0 // indirect go.uber.org/atomic v1.9.0 // indirect diff --git a/go.sum b/go.sum index ba10328498..0e77a12731 100644 --- a/go.sum +++ b/go.sum @@ -711,7 +711,6 @@ github.com/libp2p/go-yamux/v3 v3.0.1/go.mod h1:s2LsDhHbh+RfCsQoICSYt58U2f8ijtPAN github.com/libp2p/go-yamux/v3 v3.0.2/go.mod h1:s2LsDhHbh+RfCsQoICSYt58U2f8ijtPANFD8BmE74Bo= github.com/libp2p/go-yamux/v3 v3.1.2 h1:lNEy28MBk1HavUAlzKgShp+F6mn/ea1nDYWftZhFW9Q= github.com/libp2p/go-yamux/v3 v3.1.2/go.mod h1:jeLEQgLXqE2YqX1ilAClIfCMDY+0uXQUKmmb/qp0gT4= -github.com/libp2p/zeroconf/v2 v2.1.1 h1:XAuSczA96MYkVwH+LqqqCUZb2yH3krobMJ1YE+0hG2s= github.com/libp2p/zeroconf/v2 v2.1.1/go.mod h1:fuJqLnUwZTshS3U/bMRJ3+ow/v9oid1n0DmyYyNO1Xs= github.com/lightstep/lightstep-tracer-common/golang/gogo v0.0.0-20190605223551-bc2310a04743/go.mod h1:qklhhLq1aX+mtWk9cPHPzaBjWImj5ULL6C7HFJtXQMM= github.com/lightstep/lightstep-tracer-go v0.18.1/go.mod h1:jlF1pusYV4pidLvZ+XD0UBX0ZE6WURAspgAczcDHrL4= diff --git a/internal/mdns/dialable.go b/internal/mdns/dialable.go new file mode 100644 index 0000000000..1da61b7408 --- /dev/null +++ b/internal/mdns/dialable.go @@ -0,0 +1,60 @@ +// Copyright 2022 ChainSafe Systems (ON) +// SPDX-License-Identifier: LGPL-3.0-only + +package mdns + +import ( + "errors" + "fmt" + "net" + + manet "github.com/multiformats/go-multiaddr/net" +) + +var ( + ErrTCPListenAddressNotFound = errors.New("TCP listen address not found") +) + +func getMDNSIPsAndPort(network interfaceListenAddressesGetter) (ips []net.IP, port uint16) { + tcpAddresses, err := getDialableListenAddrs(network) + if err != nil { + const defaultPort = 4001 + return nil, defaultPort + } + + ips = make([]net.IP, len(tcpAddresses)) + for i := range tcpAddresses { + ips[i] = tcpAddresses[i].IP + } + port = uint16(tcpAddresses[0].Port) + + return ips, port +} + +func getDialableListenAddrs(network interfaceListenAddressesGetter) (tcpAddresses []*net.TCPAddr, err error) { + multiAddresses, err := network.InterfaceListenAddresses() + if err != nil { + return nil, fmt.Errorf("listing host interface listen addresses: %w", err) + } + + tcpAddresses = make([]*net.TCPAddr, 0, len(multiAddresses)) + for _, multiAddress := range multiAddresses { + netAddress, err := manet.ToNetAddr(multiAddress) + if err != nil { + continue + } + + tcpAddress, ok := netAddress.(*net.TCPAddr) + if !ok { + continue + } + + tcpAddresses = append(tcpAddresses, tcpAddress) + } + + if len(tcpAddresses) == 0 { + return nil, fmt.Errorf("%w: in %d multiaddresses", ErrTCPListenAddressNotFound, len(multiAddresses)) + } + + return tcpAddresses, nil +} diff --git a/internal/mdns/interfaces.go b/internal/mdns/interfaces.go new file mode 100644 index 0000000000..8021b4c579 --- /dev/null +++ b/internal/mdns/interfaces.go @@ -0,0 +1,32 @@ +// Copyright 2022 ChainSafe Systems (ON) +// SPDX-License-Identifier: LGPL-3.0-only + +package mdns + +import ( + "github.com/libp2p/go-libp2p-core/network" + "github.com/libp2p/go-libp2p-core/peer" + "github.com/multiformats/go-multiaddr" +) + +// Logger is a logger interface for the mDNS service. +type Logger interface { + Debugf(format string, args ...any) + Warnf(format string, args ...any) +} + +// IDNetworker can return the peer ID and a network interface. +type IDNetworker interface { + ID() peer.ID + Networker +} + +// Networker can return a network interface. +type Networker interface { + Network() network.Network +} + +// interfaceListenAddressesGetter returns the listen addresses of the interfaces. +type interfaceListenAddressesGetter interface { + InterfaceListenAddresses() ([]multiaddr.Multiaddr, error) +} diff --git a/internal/mdns/mdns.go b/internal/mdns/mdns.go new file mode 100644 index 0000000000..3f81fb3f25 --- /dev/null +++ b/internal/mdns/mdns.go @@ -0,0 +1,211 @@ +// Copyright 2022 ChainSafe Systems (ON) +// SPDX-License-Identifier: LGPL-3.0-only + +package mdns + +import ( + "errors" + "fmt" + "net" + "sync" + "time" + + "github.com/libp2p/go-libp2p-core/peer" + + "github.com/multiformats/go-multiaddr" + manet "github.com/multiformats/go-multiaddr/net" + "github.com/whyrusleeping/mdns" +) + +// Notifee is notified when a new peer is found. +type Notifee interface { + HandlePeerFound(peer.AddrInfo) +} + +// Service implements a mDNS service. +type Service struct { + // Dependencies and configuration injected + p2pHost IDNetworker + serviceTag string + logger Logger + notifee Notifee + + // Constant fields + pollPeriod time.Duration + + // Fields set by the Start method. + server *mdns.Server + + // Internal service management fields. + // startStopMutex is to prevent concurrent calls to Start and Stop. + startStopMutex sync.Mutex + started bool + stop chan struct{} + done chan struct{} +} + +// NewService creates and returns a new mDNS service. +func NewService(p2pHost IDNetworker, serviceTag string, + logger Logger, notifee Notifee) (service *Service) { + if serviceTag == "" { + serviceTag = "_ipfs-discovery._udp" + } + + return &Service{ + p2pHost: p2pHost, + serviceTag: serviceTag, + notifee: notifee, + logger: logger, + pollPeriod: time.Minute, + } +} + +// Start starts the mDNS service. +func (s *Service) Start() (err error) { + s.startStopMutex.Lock() + defer s.startStopMutex.Unlock() + + if s.started { + return nil + } + + ips, port := getMDNSIPsAndPort(s.p2pHost.Network()) + + hostID := s.p2pHost.ID() + + hostIDPretty := hostID.Pretty() + txt := []string{hostIDPretty} + mdnsService, err := mdns.NewMDNSService(hostIDPretty, s.serviceTag, "", "", int(port), ips, txt) + if err != nil { + return fmt.Errorf("creating mDNS service: %w", err) + } + + server, err := mdns.NewServer(&mdns.Config{Zone: mdnsService}) + if err != nil { + return fmt.Errorf("creating mDNS server: %w", err) + } + s.server = server + + s.stop = make(chan struct{}) + s.done = make(chan struct{}) + ready := make(chan struct{}) + + go s.run(ready) + // It takes a few milliseconds to launch a goroutine + // so we wait for the run goroutine to be ready. + <-ready + + s.started = true + + return nil +} + +// Stop stops the mDNS service and server. +func (s *Service) Stop() (err error) { + s.startStopMutex.Lock() + defer s.startStopMutex.Unlock() + + if !s.started { + return nil + } + + defer func() { + s.started = false + }() + close(s.stop) + <-s.done + return s.server.Shutdown() +} + +func (s *Service) run(ready chan<- struct{}) { + defer close(s.done) + + ticker := time.NewTicker(s.pollPeriod) + defer ticker.Stop() + + const queryTimeout = 5 * time.Second + params := &mdns.QueryParam{ + Domain: "local", + Service: s.serviceTag, + Timeout: queryTimeout, + } + + close(ready) + + for { + entriesListeningReady := make(chan struct{}) + entriesListeningDone := make(chan struct{}) + entriesCh := make(chan *mdns.ServiceEntry, 16) + go func() { + defer close(entriesListeningDone) + close(entriesListeningReady) + for entry := range entriesCh { + err := s.handleEntry(entry) + if err != nil { + s.logger.Warnf("handling mDNS entry: %s", err) + } + } + }() + <-entriesListeningReady + + params.Entries = entriesCh + err := mdns.Query(params) + if err != nil { + s.logger.Warnf("mdns query failed: %s", err) + } + + close(entriesCh) + <-entriesListeningDone + + select { + case <-ticker.C: + case <-s.stop: + return + } + } +} + +var ( + errEntryHasNoIP = errors.New("MDNS entry has no IP address") +) + +func (s *Service) handleEntry(entry *mdns.ServiceEntry) (err error) { + receivedPeerID, err := peer.Decode(entry.Info) + if err != nil { + return fmt.Errorf("parsing peer ID from mdns entry: %w", err) + } + + if receivedPeerID == s.p2pHost.ID() { + return nil + } + + var ip net.IP + switch { + case entry.AddrV4 != nil: + ip = entry.AddrV4 + case entry.AddrV6 != nil: + ip = entry.AddrV6 + default: + return fmt.Errorf("%w: from peer id %s", errEntryHasNoIP, receivedPeerID) + } + + tcpAddress := &net.TCPAddr{ + IP: ip, + Port: entry.Port, + } + + multiAddress, err := manet.FromNetAddr(tcpAddress) + if err != nil { + return fmt.Errorf("converting tcp address from peer id %s to multiaddress: %w", + receivedPeerID, err) + } + + addressInfo := peer.AddrInfo{ + ID: receivedPeerID, + Addrs: []multiaddr.Multiaddr{multiAddress}, + } + + s.logger.Debugf("Peer %s has addresses %s", receivedPeerID, addressInfo.Addrs) + go s.notifee.HandlePeerFound(addressInfo) + return nil +} diff --git a/internal/mdns/notifee.go b/internal/mdns/notifee.go new file mode 100644 index 0000000000..72be0c29e9 --- /dev/null +++ b/internal/mdns/notifee.go @@ -0,0 +1,42 @@ +// Copyright 2022 ChainSafe Systems (ON) +// SPDX-License-Identifier: LGPL-3.0-only + +package mdns + +import ( + "time" + + "github.com/libp2p/go-libp2p-core/peer" + "github.com/libp2p/go-libp2p-core/peerstore" + "github.com/multiformats/go-multiaddr" +) + +// AddressAdder is an interface that adds addresses. +type AddressAdder interface { + AddAddrs(p peer.ID, addrs []multiaddr.Multiaddr, ttl time.Duration) +} + +// PeerAdder adds peers. +type PeerAdder interface { + AddPeer(setID int, peerIDs ...peer.ID) +} + +// NewNotifeeTracker returns a new notifee tracker. +func NewNotifeeTracker(addressAdder AddressAdder, peerAdder PeerAdder) *NotifeeTracker { + return &NotifeeTracker{ + addressAdder: addressAdder, + peerAdder: peerAdder, + } +} + +// NotifeeTracker tracks new peers found. +type NotifeeTracker struct { + addressAdder AddressAdder + peerAdder PeerAdder +} + +// HandlePeerFound tracks the address info from the peer found. +func (n *NotifeeTracker) HandlePeerFound(p peer.AddrInfo) { + n.addressAdder.AddAddrs(p.ID, p.Addrs, peerstore.PermanentAddrTTL) + n.peerAdder.AddPeer(0, p.ID) +}