dht.go 24.4 KB
Newer Older
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
1 2
package dht

3
import (
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
4
	"bytes"
Jeromy's avatar
Jeromy committed
5
	"context"
6
	"errors"
7
	"fmt"
Aarsh Shah's avatar
Aarsh Shah committed
8
	"math"
9
	"math/rand"
10 11
	"sync"
	"time"
12

13 14 15 16 17 18
	"github.com/libp2p/go-libp2p-core/host"
	"github.com/libp2p/go-libp2p-core/network"
	"github.com/libp2p/go-libp2p-core/peer"
	"github.com/libp2p/go-libp2p-core/peerstore"
	"github.com/libp2p/go-libp2p-core/protocol"
	"github.com/libp2p/go-libp2p-core/routing"
19

20
	"github.com/libp2p/go-libp2p-kad-dht/metrics"
21
	pb "github.com/libp2p/go-libp2p-kad-dht/pb"
Aarsh Shah's avatar
Aarsh Shah committed
22
	"github.com/libp2p/go-libp2p-kad-dht/providers"
23
	"github.com/libp2p/go-libp2p-kad-dht/rtrefresh"
Aarsh Shah's avatar
Aarsh Shah committed
24 25 26
	kb "github.com/libp2p/go-libp2p-kbucket"
	record "github.com/libp2p/go-libp2p-record"
	recpb "github.com/libp2p/go-libp2p-record/pb"
27

Aarsh Shah's avatar
Aarsh Shah committed
28
	"github.com/gogo/protobuf/proto"
29 30
	ds "github.com/ipfs/go-datastore"
	logging "github.com/ipfs/go-log"
Aarsh Shah's avatar
Aarsh Shah committed
31
	"github.com/jbenet/goprocess"
Henrique Dias's avatar
Henrique Dias committed
32
	goprocessctx "github.com/jbenet/goprocess/context"
Steven Allen's avatar
Steven Allen committed
33
	"github.com/multiformats/go-base32"
34
	ma "github.com/multiformats/go-multiaddr"
Adin Schmahmann's avatar
Adin Schmahmann committed
35
	"github.com/multiformats/go-multihash"
36 37
	"go.opencensus.io/tag"
	"go.uber.org/zap"
38 39
)

40 41 42 43
var (
	logger     = logging.Logger("dht")
	baseLogger = logger.Desugar()
)
44

Steven Allen's avatar
Steven Allen committed
45 46 47 48 49 50
const (
	// BaseConnMgrScore is the base of the score set on the connection
	// manager "kbucket" tag. It is added with the common prefix length
	// between two peer IDs.
	baseConnMgrScore = 5
)
Henrique Dias's avatar
Henrique Dias committed
51

52 53 54
type mode int

const (
Steven Allen's avatar
Steven Allen committed
55 56
	modeServer mode = iota + 1
	modeClient
57 58
)

Adin Schmahmann's avatar
Adin Schmahmann committed
59 60 61 62 63
const (
	kad1 protocol.ID = "/kad/1.0.0"
	kad2 protocol.ID = "/kad/2.0.0"
)

64
const (
65 66
	kbucketTag       = "kbucket"
	protectedBuckets = 2
67 68
)

69
// IpfsDHT is an implementation of Kademlia with S/Kademlia modifications.
70
// It is used to implement the base Routing module.
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
71
type IpfsDHT struct {
72 73 74
	host      host.Host // the network services we need
	self      peer.ID   // Local peer (yourself)
	selfKey   kb.ID
75
	peerstore peerstore.Peerstore // Peer Registry
76

77
	datastore ds.Datastore // Local data
78

79
	routingTable *kb.RoutingTable // Array of routing tables for differently distanced nodes
80 81
	// ProviderManager stores & manages the provider records for this Dht peer.
	ProviderManager *providers.ProviderManager
82

83 84 85
	// manages Routing Table refresh
	rtRefreshManager *rtrefresh.RtRefreshManager

Aarsh Shah's avatar
Aarsh Shah committed
86
	birth time.Time // When this peer started up
87

88
	Validator record.Validator
89

90 91
	ctx  context.Context
	proc goprocess.Process
92 93 94

	strmap map[peer.ID]*messageSender
	smlk   sync.Mutex
95

Steven Allen's avatar
Steven Allen committed
96
	plk sync.Mutex
97

98 99
	stripedPutLocks [256]sync.Mutex

100 101
	// DHT protocols we query with. We'll only add peers to our routing
	// table if they speak these protocols.
102 103
	protocols     []protocol.ID
	protocolsStrs []string
Adin Schmahmann's avatar
Adin Schmahmann committed
104

105
	// DHT protocols we can respond to.
Adin Schmahmann's avatar
Adin Schmahmann committed
106
	serverProtocols []protocol.ID
107

108
	auto   ModeOpt
109 110 111
	mode   mode
	modeLk sync.Mutex

112
	bucketSize int
113
	alpha      int // The concurrency parameter per path
Adin Schmahmann's avatar
Adin Schmahmann committed
114
	beta       int // The number of peers closest to a target that must have responded for a query path to terminate
115

116 117 118
	queryPeerFilter        QueryFilterFunc
	routingTablePeerFilter RouteTableFilterFunc

119
	autoRefresh bool
Aarsh Shah's avatar
Aarsh Shah committed
120

121 122 123 124 125
	// A set of bootstrap peers to fallback on if all other attempts to fix
	// the routing table fail (or, e.g., this is the first time this node is
	// connecting to the network).
	bootstrapPeers []peer.AddrInfo

Aarsh Shah's avatar
Aarsh Shah committed
126
	maxRecordAge time.Duration
127

128 129 130
	// Allows disabling dht subsystems. These should _only_ be set on
	// "forked" DHTs (e.g., DHTs with custom protocols and/or private
	// networks).
131
	enableProviders, enableValues bool
132 133

	fixLowPeersChan chan struct{}
134 135
}

Matt Joiner's avatar
Matt Joiner committed
136 137 138 139
// Assert that IPFS assumptions about interfaces aren't broken. These aren't a
// guarantee, but we can use them to aid refactoring.
var (
	_ routing.ContentRouting = (*IpfsDHT)(nil)
140
	_ routing.Routing        = (*IpfsDHT)(nil)
Matt Joiner's avatar
Matt Joiner committed
141 142 143 144 145
	_ routing.PeerRouting    = (*IpfsDHT)(nil)
	_ routing.PubKeyFetcher  = (*IpfsDHT)(nil)
	_ routing.ValueStore     = (*IpfsDHT)(nil)
)

146
// New creates a new DHT with the specified host and options.
Aarsh Shah's avatar
Aarsh Shah committed
147 148 149
// Please note that being connected to a DHT peer does not necessarily imply that it's also in the DHT Routing Table.
// If the Routing Table has more than "minRTRefreshThreshold" peers, we consider a peer as a Routing Table candidate ONLY when
// we successfully get a query response from it OR if it send us a query.
150 151
func New(ctx context.Context, h host.Host, options ...Option) (*IpfsDHT, error) {
	var cfg config
Adin Schmahmann's avatar
Adin Schmahmann committed
152
	if err := cfg.apply(append([]Option{defaults}, options...)...); err != nil {
153 154
		return nil, err
	}
155 156 157
	if err := cfg.applyFallbacks(h); err != nil {
		return nil, err
	}
Aarsh Shah's avatar
Aarsh Shah committed
158

Adin Schmahmann's avatar
Adin Schmahmann committed
159 160
	if err := cfg.validate(); err != nil {
		return nil, err
161
	}
162

Aarsh Shah's avatar
Aarsh Shah committed
163 164 165 166
	dht, err := makeDHT(ctx, h, cfg)
	if err != nil {
		return nil, fmt.Errorf("failed to create DHT, err=%s", err)
	}
Adin Schmahmann's avatar
Adin Schmahmann committed
167

168
	dht.autoRefresh = cfg.routingTable.autoRefresh
169

170 171 172
	dht.maxRecordAge = cfg.maxRecordAge
	dht.enableProviders = cfg.enableProviders
	dht.enableValues = cfg.enableValues
Aarsh Shah's avatar
Aarsh Shah committed
173

174
	dht.Validator = cfg.validator
175

176
	dht.auto = cfg.mode
177
	switch cfg.mode {
178
	case ModeAuto, ModeClient:
179
		dht.mode = modeClient
180
	case ModeAutoServer, ModeServer:
181 182
		dht.mode = modeServer
	default:
183
		return nil, fmt.Errorf("invalid dht mode %d", cfg.mode)
184 185 186 187 188
	}

	if dht.mode == modeServer {
		if err := dht.moveToServerMode(); err != nil {
			return nil, err
189
		}
190
	}
191 192 193 194 195 196 197 198 199 200

	// register for event bus and network notifications
	sn, err := newSubscriberNotifiee(dht)
	if err != nil {
		return nil, err
	}
	dht.proc.Go(sn.subscribe)
	// handle providers
	dht.proc.AddChild(dht.ProviderManager.Process())

201 202 203
	if err := dht.rtRefreshManager.Start(); err != nil {
		return nil, err
	}
Aarsh Shah's avatar
Aarsh Shah committed
204 205 206 207

	// go-routine to make sure we ALWAYS have RT peer addresses in the peerstore
	// since RT membership is decoupled from connectivity
	go dht.persistRTPeersInPeerStore()
208 209 210 211

	// listens to the fix low peers chan and tries to fix the Routing Table
	dht.proc.Go(dht.fixLowPeersRoutine)

212 213
	return dht, nil
}
214

215 216 217 218
// NewDHT creates a new DHT object with the given peer as the 'local' host.
// IpfsDHT's initialized with this function will respond to DHT requests,
// whereas IpfsDHT's initialized with NewDHTClient will not.
func NewDHT(ctx context.Context, h host.Host, dstore ds.Batching) *IpfsDHT {
219
	dht, err := New(ctx, h, Datastore(dstore))
220 221 222 223 224 225 226 227 228 229
	if err != nil {
		panic(err)
	}
	return dht
}

// NewDHTClient creates a new DHT object with the given peer as the 'local'
// host. IpfsDHT clients initialized with this function will not respond to DHT
// requests. If you need a peer to respond to DHT requests, use NewDHT instead.
func NewDHTClient(ctx context.Context, h host.Host, dstore ds.Batching) *IpfsDHT {
Adin Schmahmann's avatar
Adin Schmahmann committed
230
	dht, err := New(ctx, h, Datastore(dstore), Mode(ModeClient))
231 232 233
	if err != nil {
		panic(err)
	}
234 235 236
	return dht
}

Aarsh Shah's avatar
Aarsh Shah committed
237
func makeDHT(ctx context.Context, h host.Host, cfg config) (*IpfsDHT, error) {
238
	var protocols, serverProtocols []protocol.ID
Adin Schmahmann's avatar
Adin Schmahmann committed
239 240

	// check if custom test protocols were set
241
	if cfg.v1CompatibleMode {
242 243 244 245 246 247 248 249 250 251 252 253 254 255
		// In compat mode, query/serve using the old protocol.
		//
		// DO NOT accept requests on the new protocol. Otherwise:
		// 1. We'll end up in V2 routing tables.
		// 2. We'll have V1 peers in our routing table.
		//
		// In other words, we'll pollute the V2 network.
		protocols = []protocol.ID{cfg.protocolPrefix + kad1}
		serverProtocols = []protocol.ID{cfg.protocolPrefix + kad1}
	} else {
		// In v2 mode, serve on both protocols, but only
		// query/accept peers in v2 mode.
		protocols = []protocol.ID{cfg.protocolPrefix + kad2}
		serverProtocols = []protocol.ID{cfg.protocolPrefix + kad2, cfg.protocolPrefix + kad1}
Adin Schmahmann's avatar
Adin Schmahmann committed
256 257
	}

258
	dht := &IpfsDHT{
259 260
		datastore:              cfg.datastore,
		self:                   h.ID(),
261
		selfKey:                kb.ConvertPeerID(h.ID()),
262 263 264 265 266
		peerstore:              h.Peerstore(),
		host:                   h,
		strmap:                 make(map[peer.ID]*messageSender),
		birth:                  time.Now(),
		protocols:              protocols,
267
		protocolsStrs:          protocol.ConvertToStrings(protocols),
268 269 270 271 272 273
		serverProtocols:        serverProtocols,
		bucketSize:             cfg.bucketSize,
		alpha:                  cfg.concurrency,
		beta:                   cfg.resiliency,
		queryPeerFilter:        cfg.queryPeerFilter,
		routingTablePeerFilter: cfg.routingTable.peerFilter,
274
		fixLowPeersChan:        make(chan struct{}, 1),
Jeromy's avatar
Jeromy committed
275
	}
276

277 278 279 280 281 282 283 284 285 286 287 288 289 290
	var maxLastSuccessfulOutboundThreshold time.Duration

	// The threshold is calculated based on the expected amount of time that should pass before we
	// query a peer as part of our refresh cycle.
	// To grok the Math Wizardy that produced these exact equations, please be patient as a document explaining it will
	// be published soon.
	if cfg.concurrency < cfg.bucketSize { // (alpha < K)
		l1 := math.Log(float64(1) / float64(cfg.bucketSize))                              //(Log(1/K))
		l2 := math.Log(float64(1) - (float64(cfg.concurrency) / float64(cfg.bucketSize))) // Log(1 - (alpha / K))
		maxLastSuccessfulOutboundThreshold = time.Duration(l1 / l2 * float64(cfg.routingTable.refreshInterval))
	} else {
		maxLastSuccessfulOutboundThreshold = cfg.routingTable.refreshInterval
	}

Aarsh Shah's avatar
Aarsh Shah committed
291
	// construct routing table
292 293
	// use twice the theoritical usefulness threhold to keep older peers around longer
	rt, err := makeRoutingTable(dht, cfg, 2*maxLastSuccessfulOutboundThreshold)
Aarsh Shah's avatar
Aarsh Shah committed
294 295 296 297
	if err != nil {
		return nil, fmt.Errorf("failed to construct routing table,err=%s", err)
	}
	dht.routingTable = rt
298
	dht.bootstrapPeers = cfg.bootstrapPeers
Aarsh Shah's avatar
Aarsh Shah committed
299

300 301 302 303 304 305 306
	// rt refresh manager
	rtRefresh, err := makeRtRefreshManager(dht, cfg, maxLastSuccessfulOutboundThreshold)
	if err != nil {
		return nil, fmt.Errorf("failed to construct RT Refresh Manager,err=%s", err)
	}
	dht.rtRefreshManager = rtRefresh

307
	// create a DHT proc with the given context
308 309 310
	dht.proc = goprocessctx.WithContextAndTeardown(ctx, func() error {
		return rtRefresh.Close()
	})
Aarsh Shah's avatar
Aarsh Shah committed
311 312 313 314 315 316

	// create a tagged context derived from the original context
	ctxTags := dht.newContextWithLocalTags(ctx)
	// the DHT context should be done when the process is closed
	dht.ctx = goprocessctx.WithProcessClosing(ctxTags, dht.proc)

Alan Shaw's avatar
Alan Shaw committed
317 318 319 320 321
	pm, err := providers.NewProviderManager(dht.ctx, h.ID(), cfg.datastore, cfg.providersOptions...)
	if err != nil {
		return nil, err
	}
	dht.ProviderManager = pm
322

Aarsh Shah's avatar
Aarsh Shah committed
323
	return dht, nil
Jeromy's avatar
Jeromy committed
324 325
}

326 327 328 329 330
func makeRtRefreshManager(dht *IpfsDHT, cfg config, maxLastSuccessfulOutboundThreshold time.Duration) (*rtrefresh.RtRefreshManager, error) {
	keyGenFnc := func(cpl uint) (string, error) {
		p, err := dht.routingTable.GenRandPeerID(cpl)
		return string(p), err
	}
Aarsh Shah's avatar
Aarsh Shah committed
331

332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348
	queryFnc := func(ctx context.Context, key string) error {
		_, err := dht.GetClosestPeers(ctx, key)
		return err
	}

	r, err := rtrefresh.NewRtRefreshManager(
		dht.host, dht.routingTable, cfg.routingTable.autoRefresh,
		keyGenFnc,
		queryFnc,
		cfg.routingTable.refreshQueryTimeout,
		cfg.routingTable.refreshInterval,
		maxLastSuccessfulOutboundThreshold)

	return r, err
}

func makeRoutingTable(dht *IpfsDHT, cfg config, maxLastSuccessfulOutboundThreshold time.Duration) (*kb.RoutingTable, error) {
349
	rt, err := kb.NewRoutingTable(cfg.bucketSize, dht.selfKey, time.Minute, dht.host.Peerstore(), maxLastSuccessfulOutboundThreshold)
Aarsh Shah's avatar
Aarsh Shah committed
350
	cmgr := dht.host.ConnManager()
Aarsh Shah's avatar
Aarsh Shah committed
351 352

	rt.PeerAdded = func(p peer.ID) {
353
		commonPrefixLen := kb.CommonPrefixLen(dht.selfKey, kb.ConvertPeerID(p))
354 355 356 357 358
		if commonPrefixLen < protectedBuckets {
			cmgr.Protect(p, kbucketTag)
		} else {
			cmgr.TagPeer(p, kbucketTag, baseConnMgrScore)
		}
Aarsh Shah's avatar
Aarsh Shah committed
359 360
	}
	rt.PeerRemoved = func(p peer.ID) {
361
		cmgr.Unprotect(p, kbucketTag)
362
		cmgr.UntagPeer(p, kbucketTag)
363 364 365

		// try to fix the RT
		dht.fixRTIfNeeded()
Aarsh Shah's avatar
Aarsh Shah committed
366 367 368 369
	}

	return rt, err
}
370

Will Scott's avatar
Will Scott committed
371 372 373 374 375
// Mode allows introspection of the operation mode of the DHT
func (dht *IpfsDHT) Mode() ModeOpt {
	return dht.auto
}

376
// fixLowPeersRoutine tries to get more peers into the routing table if we're below the threshold
377
func (dht *IpfsDHT) fixLowPeersRoutine(proc goprocess.Process) {
378 379
	ticker := time.NewTicker(periodicBootstrapInterval)
	defer ticker.Stop()
380

381 382 383
	for {
		select {
		case <-dht.fixLowPeersChan:
384
		case <-ticker.C:
385 386 387
		case <-proc.Closing():
			return
		}
388

389 390 391 392
		if dht.routingTable.Size() > minRTRefreshThreshold {
			continue
		}

393 394
		// we try to add all peers we are connected to to the Routing Table
		// in case they aren't already there.
395 396 397 398
		for _, p := range dht.host.Network().Peers() {
			dht.peerFound(dht.Context(), p, false)
		}

399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439
		// TODO Active Bootstrapping
		// We should first use non-bootstrap peers we knew of from previous
		// snapshots of the Routing Table before we connect to the bootstrappers.
		// See https://github.com/libp2p/go-libp2p-kad-dht/issues/387.
		if dht.routingTable.Size() == 0 {
			if len(dht.bootstrapPeers) == 0 {
				// No point in continuing, we have no peers!
				continue
			}

			found := 0
			for _, i := range rand.Perm(len(dht.bootstrapPeers)) {
				ai := dht.bootstrapPeers[i]
				err := dht.Host().Connect(dht.Context(), ai)
				if err == nil {
					found++
				} else {
					logger.Warnw("failed to bootstrap", "peer", ai.ID, "error", err)
				}

				// Wait for two bootstrap peers, or try them all.
				//
				// Why two? In theory, one should be enough
				// normally. However, if the network were to
				// restart and everyone connected to just one
				// bootstrapper, we'll end up with a mostly
				// partitioned network.
				//
				// So we always bootstrap with two random peers.
				if found == maxNBoostrappers {
					break
				}
			}
		}

		// if we still don't have peers in our routing table(probably because Identify hasn't completed),
		// there is no point in triggering a Refresh.
		if dht.routingTable.Size() == 0 {
			continue
		}

440
		if dht.autoRefresh {
441
			dht.rtRefreshManager.RefreshNoWait()
442 443 444 445 446
		}
	}

}

Aarsh Shah's avatar
Aarsh Shah committed
447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465
// TODO This is hacky, horrible and the programmer needs to have his mother called a hamster.
// SHOULD be removed once https://github.com/libp2p/go-libp2p/issues/800 goes in.
func (dht *IpfsDHT) persistRTPeersInPeerStore() {
	tickr := time.NewTicker(peerstore.RecentlyConnectedAddrTTL / 3)
	defer tickr.Stop()

	for {
		select {
		case <-tickr.C:
			ps := dht.routingTable.ListPeers()
			for _, p := range ps {
				dht.peerstore.UpdateAddrs(p, peerstore.RecentlyConnectedAddrTTL, peerstore.RecentlyConnectedAddrTTL)
			}
		case <-dht.ctx.Done():
			return
		}
	}
}

Jeromy's avatar
Jeromy committed
466
// putValueToPeer stores the given key/value pair at the peer 'p'
467 468
func (dht *IpfsDHT) putValueToPeer(ctx context.Context, p peer.ID, rec *recpb.Record) error {
	pmes := pb.NewMessage(pb.Message_PUT_VALUE, rec.Key, 0)
469
	pmes.Record = rec
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
470
	rpmes, err := dht.sendRequest(ctx, p, pmes)
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
471
	if err != nil {
Adin Schmahmann's avatar
Adin Schmahmann committed
472
		logger.Debugw("failed to put value to peer", "to", p, "key", loggableRecordKeyBytes(rec.Key), "error", err)
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
473 474
		return err
	}
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
475

476
	if !bytes.Equal(rpmes.GetRecord().Value, pmes.GetRecord().Value) {
Steven Allen's avatar
Steven Allen committed
477
		logger.Infow("value not put correctly", "put-message", pmes, "get-message", rpmes)
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
478 479
		return errors.New("value not put correctly")
	}
gpestana's avatar
gpestana committed
480

Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
481
	return nil
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
482 483
}

484 485
var errInvalidRecord = errors.New("received invalid record")

486 487
// getValueOrPeers queries a particular peer p for the value for
// key. It returns either the value or a list of closer peers.
488
// NOTE: It will update the dht's peerstore with any new addresses
489
// it finds for the given peer.
490
func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.ID, key string) (*recpb.Record, []*peer.AddrInfo, error) {
491
	pmes, err := dht.getValueSingle(ctx, p, key)
492
	if err != nil {
493
		return nil, nil, err
494 495
	}

496 497 498
	// Perhaps we were given closer peers
	peers := pb.PBPeersToPeerInfos(pmes.GetCloserPeers())

499
	if record := pmes.GetRecord(); record != nil {
500
		// Success! We were given the value
Steven Allen's avatar
Steven Allen committed
501
		logger.Debug("got value")
502

503
		// make sure record is valid.
504
		err = dht.Validator.Validate(string(record.GetKey()), record.GetValue())
505
		if err != nil {
Steven Allen's avatar
Steven Allen committed
506
			logger.Debug("received invalid record (discarded)")
507 508
			// return a sentinal to signify an invalid record was received
			err = errInvalidRecord
George Antoniadis's avatar
George Antoniadis committed
509
			record = new(recpb.Record)
510
		}
511
		return record, peers, err
512
	}
513

514 515 516 517
	if len(peers) > 0 {
		return nil, peers, nil
	}

518
	return nil, nil, routing.ErrNotFound
519 520
}

521
// getValueSingle simply performs the get value RPC with the given parameters
522
func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.ID, key string) (*pb.Message, error) {
523
	pmes := pb.NewMessage(pb.Message_GET_VALUE, []byte(key), 0)
Steven Allen's avatar
Steven Allen committed
524
	return dht.sendRequest(ctx, p, pmes)
Jeromy's avatar
Jeromy committed
525 526
}

527
// getLocal attempts to retrieve the value from the datastore
528
func (dht *IpfsDHT) getLocal(key string) (*recpb.Record, error) {
Adin Schmahmann's avatar
Adin Schmahmann committed
529
	logger.Debugw("finding value in datastore", "key", loggableRecordKeyString(key))
Steven Allen's avatar
Steven Allen committed
530

531
	rec, err := dht.getRecordFromDatastore(mkDsKey(key))
532
	if err != nil {
Adin Schmahmann's avatar
Adin Schmahmann committed
533
		logger.Warnw("get local failed", "key", loggableRecordKeyString(key), "error", err)
534 535
		return nil, err
	}
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
536

537
	// Double check the key. Can't hurt.
538
	if rec != nil && string(rec.GetKey()) != key {
Adin Schmahmann's avatar
Adin Schmahmann committed
539
		logger.Errorw("BUG: found a DHT record that didn't match it's key", "expected", loggableRecordKeyString(key), "got", rec.GetKey())
Steven Allen's avatar
Steven Allen committed
540
		return nil, nil
541 542

	}
543
	return rec, nil
544 545
}

546
// putLocal stores the key value pair in the datastore
547
func (dht *IpfsDHT) putLocal(key string, rec *recpb.Record) error {
548 549
	data, err := proto.Marshal(rec)
	if err != nil {
Adin Schmahmann's avatar
Adin Schmahmann committed
550
		logger.Warnw("failed to put marshal record for local put", "error", err, "key", loggableRecordKeyString(key))
551 552 553
		return err
	}

554
	return dht.datastore.Put(mkDsKey(key), data)
555
}
556

Aarsh Shah's avatar
Aarsh Shah committed
557
// peerFound signals the routingTable that we've found a peer that
Aarsh Shah's avatar
Aarsh Shah committed
558
// might support the DHT protocol.
Aarsh Shah's avatar
Aarsh Shah committed
559 560
// If we have a connection a peer but no exchange of a query RPC ->
//    LastQueriedAt=time.Now (so we don't ping it for some time for a liveliness check)
561 562
//    LastUsefulAt=0
// If we connect to a peer and then exchange a query RPC ->
Aarsh Shah's avatar
Aarsh Shah committed
563 564 565 566 567 568 569
//    LastQueriedAt=time.Now (same reason as above)
//    LastUsefulAt=time.Now (so we give it some life in the RT without immediately evicting it)
// If we query a peer we already have in our Routing Table ->
//    LastQueriedAt=time.Now()
//    LastUsefulAt remains unchanged
// If we connect to a peer we already have in the RT but do not exchange a query (rare)
//    Do Nothing.
Aarsh Shah's avatar
Aarsh Shah committed
570
func (dht *IpfsDHT) peerFound(ctx context.Context, p peer.ID, queryPeer bool) {
571 572 573
	if c := baseLogger.Check(zap.DebugLevel, "peer found"); c != nil {
		c.Write(zap.String("peer", p.String()))
	}
Aarsh Shah's avatar
Aarsh Shah committed
574 575
	b, err := dht.validRTPeer(p)
	if err != nil {
Steven Allen's avatar
Steven Allen committed
576
		logger.Errorw("failed to validate if peer is a DHT peer", "peer", p, "error", err)
Aarsh Shah's avatar
Aarsh Shah committed
577
	} else if b {
Aarsh Shah's avatar
Aarsh Shah committed
578
		newlyAdded, err := dht.routingTable.TryAddPeer(p, queryPeer)
Steven Allen's avatar
Steven Allen committed
579 580 581 582
		if err != nil {
			// peer not added.
			return
		}
583
		if !newlyAdded && queryPeer {
Aarsh Shah's avatar
Aarsh Shah committed
584 585 586
			// the peer is already in our RT, but we just successfully queried it and so let's give it a
			// bump on the query time so we don't ping it too soon for a liveliness check.
			dht.routingTable.UpdateLastSuccessfulOutboundQueryAt(p, time.Now())
Aarsh Shah's avatar
Aarsh Shah committed
587 588
		}
	}
Aarsh Shah's avatar
Aarsh Shah committed
589 590
}

Aarsh Shah's avatar
Aarsh Shah committed
591
// peerStoppedDHT signals the routing table that a peer is unable to responsd to DHT queries anymore.
Aarsh Shah's avatar
Aarsh Shah committed
592
func (dht *IpfsDHT) peerStoppedDHT(ctx context.Context, p peer.ID) {
Steven Allen's avatar
Steven Allen committed
593
	logger.Debugw("peer stopped dht", "peer", p)
Aarsh Shah's avatar
Aarsh Shah committed
594 595
	// A peer that does not support the DHT protocol is dead for us.
	// There's no point in talking to anymore till it starts supporting the DHT protocol again.
Aarsh Shah's avatar
Aarsh Shah committed
596
	dht.routingTable.RemovePeer(p)
Aarsh Shah's avatar
Aarsh Shah committed
597 598
}

599 600 601 602
func (dht *IpfsDHT) fixRTIfNeeded() {
	select {
	case dht.fixLowPeersChan <- struct{}{}:
	default:
Aarsh Shah's avatar
Aarsh Shah committed
603
	}
604
}
Jeromy's avatar
Jeromy committed
605

Jeromy's avatar
Jeromy committed
606
// FindLocal looks for a peer with a given ID connected to this dht and returns the peer and the table it was found in.
607
func (dht *IpfsDHT) FindLocal(id peer.ID) peer.AddrInfo {
608
	switch dht.host.Network().Connectedness(id) {
609
	case network.Connected, network.CanConnect:
610 611
		return dht.peerstore.PeerInfo(id)
	default:
612
		return peer.AddrInfo{}
Jeromy's avatar
Jeromy committed
613 614
	}
}
615

Jeromy's avatar
Jeromy committed
616
// findPeerSingle asks peer 'p' if they know where the peer with id 'id' is
617
func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.ID, id peer.ID) (*pb.Message, error) {
618
	pmes := pb.NewMessage(pb.Message_FIND_NODE, []byte(id), 0)
Steven Allen's avatar
Steven Allen committed
619
	return dht.sendRequest(ctx, p, pmes)
620
}
621

Adin Schmahmann's avatar
Adin Schmahmann committed
622 623
func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.ID, key multihash.Multihash) (*pb.Message, error) {
	pmes := pb.NewMessage(pb.Message_GET_PROVIDERS, key, 0)
Steven Allen's avatar
Steven Allen committed
624
	return dht.sendRequest(ctx, p, pmes)
Jeromy's avatar
Jeromy committed
625 626
}

627
// nearestPeersToQuery returns the routing tables closest peers.
628
func (dht *IpfsDHT) nearestPeersToQuery(pmes *pb.Message, count int) []peer.ID {
629
	closer := dht.routingTable.NearestPeers(kb.ConvertKey(string(pmes.GetKey())), count)
630 631 632
	return closer
}

Aarsh Shah's avatar
Aarsh Shah committed
633
// betterPeersToQuery returns nearestPeersToQuery with some additional filtering
634
func (dht *IpfsDHT) betterPeersToQuery(pmes *pb.Message, from peer.ID, count int) []peer.ID {
635
	closer := dht.nearestPeersToQuery(pmes, count)
636 637 638

	// no node? nil
	if closer == nil {
Steven Allen's avatar
Steven Allen committed
639
		logger.Infow("no closer peers to send", from)
640 641 642
		return nil
	}

Steven Allen's avatar
Steven Allen committed
643
	filtered := make([]peer.ID, 0, len(closer))
Jeromy's avatar
Jeromy committed
644 645 646
	for _, clp := range closer {

		// == to self? thats bad
Jeromy's avatar
Jeromy committed
647
		if clp == dht.self {
Matt Joiner's avatar
Matt Joiner committed
648
			logger.Error("BUG betterPeersToQuery: attempted to return self! this shouldn't happen...")
649 650
			return nil
		}
651
		// Dont send a peer back themselves
652
		if clp == from {
653 654 655
			continue
		}

Jeromy's avatar
Jeromy committed
656
		filtered = append(filtered, clp)
657 658
	}

659 660
	// ok seems like closer nodes
	return filtered
661 662
}

663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680
func (dht *IpfsDHT) setMode(m mode) error {
	dht.modeLk.Lock()
	defer dht.modeLk.Unlock()

	if m == dht.mode {
		return nil
	}

	switch m {
	case modeServer:
		return dht.moveToServerMode()
	case modeClient:
		return dht.moveToClientMode()
	default:
		return fmt.Errorf("unrecognized dht mode: %d", m)
	}
}

Adin Schmahmann's avatar
Adin Schmahmann committed
681 682 683
// moveToServerMode advertises (via libp2p identify updates) that we are able to respond to DHT queries and sets the appropriate stream handlers.
// Note: We may support responding to queries with protocols aside from our primary ones in order to support
// interoperability with older versions of the DHT protocol.
684 685
func (dht *IpfsDHT) moveToServerMode() error {
	dht.mode = modeServer
Adin Schmahmann's avatar
Adin Schmahmann committed
686
	for _, p := range dht.serverProtocols {
687 688 689 690 691
		dht.host.SetStreamHandler(p, dht.handleNewStream)
	}
	return nil
}

Adin Schmahmann's avatar
Adin Schmahmann committed
692 693 694 695 696
// moveToClientMode stops advertising (and rescinds advertisements via libp2p identify updates) that we are able to
// respond to DHT queries and removes the appropriate stream handlers. We also kill all inbound streams that were
// utilizing the handled protocols.
// Note: We may support responding to queries with protocols aside from our primary ones in order to support
// interoperability with older versions of the DHT protocol.
697 698
func (dht *IpfsDHT) moveToClientMode() error {
	dht.mode = modeClient
Adin Schmahmann's avatar
Adin Schmahmann committed
699
	for _, p := range dht.serverProtocols {
700 701 702 703
		dht.host.RemoveStreamHandler(p)
	}

	pset := make(map[protocol.ID]bool)
Adin Schmahmann's avatar
Adin Schmahmann committed
704
	for _, p := range dht.serverProtocols {
705 706 707 708 709 710 711
		pset[p] = true
	}

	for _, c := range dht.host.Network().Conns() {
		for _, s := range c.GetStreams() {
			if pset[s.Protocol()] {
				if s.Stat().Direction == network.DirInbound {
Steven Allen's avatar
Steven Allen committed
712
					_ = s.Reset()
713 714 715 716 717 718 719 720 721 722 723 724 725
				}
			}
		}
	}
	return nil
}

func (dht *IpfsDHT) getMode() mode {
	dht.modeLk.Lock()
	defer dht.modeLk.Unlock()
	return dht.mode
}

Alan Shaw's avatar
Alan Shaw committed
726
// Context returns the DHT's context.
727 728 729 730
func (dht *IpfsDHT) Context() context.Context {
	return dht.ctx
}

Alan Shaw's avatar
Alan Shaw committed
731
// Process returns the DHT's process.
732 733 734 735
func (dht *IpfsDHT) Process() goprocess.Process {
	return dht.proc
}

Alan Shaw's avatar
Alan Shaw committed
736
// RoutingTable returns the DHT's routingTable.
ZhengQi's avatar
ZhengQi committed
737 738 739 740
func (dht *IpfsDHT) RoutingTable() *kb.RoutingTable {
	return dht.routingTable
}

Alan Shaw's avatar
Alan Shaw committed
741
// Close calls Process Close.
742 743 744
func (dht *IpfsDHT) Close() error {
	return dht.proc.Close()
}
745 746 747 748

func mkDsKey(s string) ds.Key {
	return ds.NewKey(base32.RawStdEncoding.EncodeToString([]byte(s)))
}
749

Alan Shaw's avatar
Alan Shaw committed
750
// PeerID returns the DHT node's Peer ID.
751 752 753 754
func (dht *IpfsDHT) PeerID() peer.ID {
	return dht.self
}

Alan Shaw's avatar
Alan Shaw committed
755
// PeerKey returns a DHT key, converted from the DHT node's Peer ID.
756 757 758 759
func (dht *IpfsDHT) PeerKey() []byte {
	return kb.ConvertPeerID(dht.self)
}

Alan Shaw's avatar
Alan Shaw committed
760
// Host returns the libp2p host this DHT is operating with.
761 762 763 764
func (dht *IpfsDHT) Host() host.Host {
	return dht.host
}

Alan Shaw's avatar
Alan Shaw committed
765
// Ping sends a ping message to the passed peer and waits for a response.
766 767 768 769
func (dht *IpfsDHT) Ping(ctx context.Context, p peer.ID) error {
	req := pb.NewMessage(pb.Message_PING, nil, 0)
	resp, err := dht.sendRequest(ctx, p, req)
	if err != nil {
Steven Allen's avatar
Steven Allen committed
770
		return fmt.Errorf("sending request: %w", err)
771 772
	}
	if resp.Type != pb.Message_PING {
Steven Allen's avatar
Steven Allen committed
773
		return fmt.Errorf("got unexpected response type: %v", resp.Type)
774 775 776
	}
	return nil
}
777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792

// newContextWithLocalTags returns a new context.Context with the InstanceID and
// PeerID keys populated. It will also take any extra tags that need adding to
// the context as tag.Mutators.
func (dht *IpfsDHT) newContextWithLocalTags(ctx context.Context, extraTags ...tag.Mutator) context.Context {
	extraTags = append(
		extraTags,
		tag.Upsert(metrics.KeyPeerID, dht.self.Pretty()),
		tag.Upsert(metrics.KeyInstanceID, fmt.Sprintf("%p", dht)),
	)
	ctx, _ = tag.New(
		ctx,
		extraTags...,
	) // ignoring error as it is unrelated to the actual function of this code.
	return ctx
}
793 794 795 796 797 798 799 800

func (dht *IpfsDHT) maybeAddAddrs(p peer.ID, addrs []ma.Multiaddr, ttl time.Duration) {
	// Don't add addresses for self or our connected peers. We have better ones.
	if p == dht.self || dht.host.Network().Connectedness(p) == network.Connected {
		return
	}
	dht.peerstore.AddAddrs(p, addrs, ttl)
}