Skip to content

Commit 54b6682

Browse files
authored
fix(service-disco): service discovery key bucketing (#2592)
1 parent d3c396b commit 54b6682

8 files changed

Lines changed: 160 additions & 190 deletions

File tree

libp2p/protocols/kademlia.nim

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,8 @@ proc refreshTable*(
3434
if not (forceRefresh or bucket.isStale()):
3535
continue
3636

37-
let randomKey = randomKeyInBucket(rtable.selfId, i, kad.rng)
37+
let randomKey =
38+
randomKeyInBucket(rtable.selfId, i, kad.rng, rtable.config.maxBuckets)
3839
discard await kad.findNode(randomKey, rtable)
3940

4041
proc bootstrap*(

libp2p/protocols/kademlia/routing_table.nim

Lines changed: 39 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
# SPDX-License-Identifier: Apache-2.0 OR MIT
22
# Copyright (c) Status Research & Development GmbH
33

4-
import algorithm, sequtils
4+
import algorithm, sequtils, math, bitops
55
import chronos, chronicles, results
66
import ./types
77
import ./kademlia_metrics
@@ -40,17 +40,20 @@ proc new*(
4040
selfId: selfId, localNodeId: localNodeId.get(selfId), buckets: @[], config: config
4141
)
4242

43-
proc bucketIndex*(
44-
selfId, key: Key, hasher: Opt[XorDHasher], selfIdPreHashed = false
45-
): int =
43+
proc bucketIndex*(rtable: RoutingTable, key: Key): int =
4644
let
4745
selfHash =
48-
if selfIdPreHashed:
49-
selfId
46+
if rtable.config.selfIdPreHashed:
47+
rtable.selfId
5048
else:
51-
selfId.hashFor(hasher)
52-
keyHash = key.hashFor(hasher)
53-
return xorDistance(selfHash, keyHash).leadingZeros
49+
rtable.selfId.hashFor(rtable.config.hasher)
50+
keyHash = key.hashFor(rtable.config.hasher)
51+
lz = xorDistance(selfHash, keyHash).leadingZeros
52+
53+
if rtable.config.maxBuckets <= 1:
54+
return 0
55+
56+
return min(((lz * rtable.config.maxBuckets) div 256), rtable.config.maxBuckets - 1)
5457

5558
proc peerIndexInBucket(bucket: Bucket, nodeId: Key): Opt[int] =
5659
for i, p in bucket.peers:
@@ -97,13 +100,7 @@ proc insert*(rtable: RoutingTable, nodeId: Key): bool =
97100
debug "Cannot insert self in routing table", nodeId = nodeId
98101
return false # No self insertion
99102

100-
let idx = bucketIndex(
101-
rtable.selfId, nodeId, rtable.config.hasher, rtable.config.selfIdPreHashed
102-
)
103-
if idx >= rtable.config.maxBuckets:
104-
debug "Cannot insert node, max buckets have been reached",
105-
nodeId = nodeId, bucketIdx = idx, maxBuckets = rtable.config.maxBuckets
106-
return false
103+
let idx = rtable.bucketIndex(nodeId)
107104

108105
if idx >= rtable.buckets.len:
109106
# expand buckets lazily if needed
@@ -208,32 +205,34 @@ proc isStale*(bucket: Bucket): bool =
208205
return true
209206
return false
210207

211-
proc randomKeyInBucket*(selfId: Key, bucketIndex: int, rng: Rng): Key =
212-
var raw = selfId
213-
214-
# zero out higher bits
215-
for i in 0 ..< bucketIndex:
216-
let byteIdx = i div 8
217-
let bitInByte = 7 - (i mod 8)
218-
raw[byteIdx] = raw[byteIdx] and not (1'u8 shl bitInByte)
219-
220-
# flip the target bit
221-
let tgtByte = bucketIndex div 8
222-
let tgtBitInByte = 7 - (bucketIndex mod 8)
223-
raw[tgtByte] = raw[tgtByte] xor (1'u8 shl tgtBitInByte)
224-
225-
# randomize lower bits of the boundary byte
226-
let lsbMask = (1'u8 shl tgtBitInByte) - 1
227-
if lsbMask != 0:
228-
var rb: array[1, byte]
229-
rng.generate(rb)
230-
raw[tgtByte] = (raw[tgtByte] and not lsbMask) or (rb[0] and lsbMask)
208+
proc randomKeyInBucket*(
209+
selfId: Key, bucketIndex: int, rng: Rng, maxBuckets: int = DefaultMaxBuckets
210+
): Key =
211+
let
212+
index = clamp(bucketIndex, 0, (maxBuckets - 1))
213+
buckets = max(1, maxBuckets)
214+
leadingZeros = index * 256 div buckets
215+
(byteIdx, rem) = divmod(leadingZeros, 8)
216+
boundBitIdx = 7 - rem
217+
218+
var key = selfId
219+
220+
# For the boundary byte, from 0 to boundBitIdx the bits should be random.
221+
for i in byteIdx ..< IdLength:
222+
rng.generate(key[i])
223+
224+
# For the boundary byte, from boundBitIdx to 7 the bits should be the same as seflId.
225+
for i in boundBitIdx .. 7:
226+
if selfId[byteIdx].testBit(i):
227+
key[byteIdx].setBit(i)
228+
else:
229+
key[byteIdx].clearBit(i)
231230

232-
# randomize remaining bytes
233-
if tgtByte + 1 < raw.len:
234-
rng.generate(raw.toOpenArray(tgtByte + 1, raw.len - 1))
231+
# The bit at the boundary should be the opposite of selfId. So that the XOR is not 0.
232+
if selfId[byteIdx].testBit(boundBitIdx) == key[byteIdx].testBit(boundBitIdx):
233+
key[byteIdx].flipBit(boundBitIdx)
235234

236-
return raw
235+
key
237236

238237
proc allKeys*(bucket: Bucket): seq[Key] =
239238
return bucket.peers.mapIt(it.nodeId)

tests/libp2p/kademlia/test_find.nim

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -205,9 +205,7 @@ suite "KadDHT Find":
205205
check not kads[0].hasKey(kads[2].rtable.selfId)
206206

207207
# Make kads[1]'s bucket stale to trigger refresh
208-
let bucketIdx = bucketIndex(
209-
kads[0].rtable.selfId, kads[1].rtable.selfId, kads[0].rtable.config.hasher
210-
)
208+
let bucketIdx = kads[0].rtable.bucketIndex(kads[1].rtable.selfId)
211209
makeBucketStale(kads[0].rtable.buckets[bucketIdx])
212210

213211
check kads[0].rtable.buckets[bucketIdx].isStale()

tests/libp2p/kademlia/test_routing_table.nim

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ suite "KadDHT Routing Table":
2121
let other = testKey(0b10000000)
2222
discard rt.insert(other)
2323

24-
let idx = bucketIndex(selfId, other, Opt.none(XorDHasher))
24+
let idx = rt.bucketIndex(other)
2525
check:
2626
rt.buckets.len > idx
2727
rt.buckets[idx].peers.len == 1
@@ -161,8 +161,10 @@ suite "KadDHT Routing Table":
161161

162162
test "randomKeyInBucket returns id at correct distance":
163163
let selfId = testKey(0)
164+
var rt =
165+
RoutingTable.new(selfId, RoutingTableConfig.new(hasher = Opt.some(noOpHasher)))
164166
var rid = randomKeyInBucket(selfId, TargetBucket, rng())
165-
let idx = bucketIndex(selfId, rid, Opt.some(noOpHasher))
167+
let idx = rt.bucketIndex(rid)
166168
check:
167169
idx == TargetBucket
168170
rid != selfId

tests/libp2p/service_discovery/component/test_advertise_discover.nim

Lines changed: 0 additions & 62 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@ import
77
../../../../libp2p/[
88
crypto/crypto,
99
peerid,
10-
protocols/kademlia/routing_table,
1110
protocols/kademlia/types,
1211
protocols/service_discovery,
1312
protocols/service_discovery/advertiser,
@@ -357,67 +356,6 @@ suite "Service Discovery Component - Advertise Discover":
357356
let foundB = await discovererNode.lookup(serviceBId)
358357
foundA.containsPeer(advertiserNode) and foundB.containsPeer(advertiserNode)
359358

360-
asyncTest "addProvidedService drops a known registrar at Kad bucketIndex 16":
361-
# TODO: nim-libp2p#2499 service-disco: valid service-table peers are dropped when Kad bucketIndex is 16 or higher
362-
# Precomputed private key to pin the registrar identity and service name.
363-
# Their service table bucket index is 18 for this serviceId.
364-
# A ServiceDiscovery service table accepts indexes 0 through 15.
365-
const
366-
serviceName = "service-bucket-indexing-component"
367-
droppedRegistrarPrivateKey =
368-
"080112404911f67066496ac09c4924a8ef86ec0d8532c9fb50046aa27fdb4132" &
369-
"130a2c8ecf17e8057d0e5c6fcef5a3aeec98d70b43feffac4301c141a08071498ce89a1e"
370-
371-
let
372-
conf = ServiceDiscoveryConfig.new(safetyParam = 0.0)
373-
droppedRegistrarKey = PrivateKey.init(droppedRegistrarPrivateKey).get()
374-
droppedRegistrarNode = setupServiceDiscoveryNode(
375-
discoConfig = conf, privateKey = Opt.some(droppedRegistrarKey)
376-
)
377-
workingRegistrarNode = setupServiceDiscoveryNode(discoConfig = conf)
378-
advertiserNode = setupServiceDiscoveryNode(discoConfig = conf)
379-
380-
startAndDeferStop(@[droppedRegistrarNode, workingRegistrarNode, advertiserNode])
381-
await connect(droppedRegistrarNode, advertiserNode)
382-
await connect(workingRegistrarNode, advertiserNode)
383-
384-
# The advertiser knows both registrars in the main Kad routing table.
385-
# The droppedRegistrar sits just outside the service table bucket range.
386-
# The workingRegistrar sits inside the service table bucket range.
387-
let
388-
service = makeServiceInfo(serviceName)
389-
serviceId = service.id.hashServiceId()
390-
droppedRegistrarPeerKey = droppedRegistrarNode.switch.peerInfo.peerId.toKey()
391-
workingRegistrarPeerKey = workingRegistrarNode.switch.peerInfo.peerId.toKey()
392-
393-
check:
394-
bucketIndex(
395-
serviceId, droppedRegistrarPeerKey, Opt.none(XorDHasher), selfIdPreHashed = true
396-
) > conf.bucketsCount
397-
bucketIndex(
398-
serviceId, workingRegistrarPeerKey, Opt.none(XorDHasher), selfIdPreHashed = true
399-
) < conf.bucketsCount
400-
advertiserNode.rtable.hasPeer(droppedRegistrarPeerKey)
401-
advertiserNode.rtable.hasPeer(workingRegistrarPeerKey)
402-
403-
# addProvidedService seeds the service table from the main Kad routing table.
404-
advertiserNode.addProvidedService(service)
405-
406-
# The in-range registrar is kept and receives the ad.
407-
# The out-of-range registrar is absent from the service table.
408-
# Only the remote task is tracked in running (local registration uses a dedicated loop).
409-
let serviceTable = advertiserNode.rtManager.getTable(serviceId).get()
410-
check:
411-
serviceTable.hasPeer(workingRegistrarPeerKey)
412-
not serviceTable.hasPeer(droppedRegistrarPeerKey)
413-
advertiserNode.advertiser.running.len() == 1
414-
# only the remote task (local registration is a separate dedicated loop)
415-
droppedRegistrarNode.countAdsInCache(serviceId) == 0
416-
417-
checkUntilTimeout:
418-
workingRegistrarNode.countAdsInCache(serviceId) == 1
419-
droppedRegistrarNode.countAdsInCache(serviceId) == 0
420-
421359
asyncTest "advertiser registers with peers discovered after addProvidedService":
422360
let conf = ServiceDiscoveryConfig.new(safetyParam = 0.0)
423361
let initialRegistrar = setupServiceDiscoveryNode(discoConfig = conf)

tests/libp2p/service_discovery/component/test_lookup_get_ads.nim

Lines changed: 15 additions & 71 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@ import
77
../../../../libp2p/[
88
crypto/crypto,
99
peerid,
10-
protocols/kademlia/routing_table,
1110
protocols/kademlia/types,
1211
protocols/service_discovery/advertiser,
1312
protocols/service_discovery/discoverer,
@@ -19,56 +18,6 @@ import ../../../../libp2p/protocols/kademlia/protobuf as kad_protobuf
1918
import ../../../tools/[lifecycle, unittest]
2019
import ../utils
2120

22-
const MaxRegistrarSetupAttempts = 1000
23-
24-
proc bucketOf(r: ServiceDiscovery, serviceId: ServiceId): int =
25-
bucketIndex(
26-
serviceId,
27-
r.switch.peerInfo.peerId.toKey(),
28-
Opt.none(XorDHasher),
29-
selfIdPreHashed = true,
30-
)
31-
32-
proc setupRegistrarsInDistinctBuckets(
33-
conf: ServiceDiscoveryConfig, serviceId: ServiceId
34-
): tuple[queriedFirst, queriedSecond: ServiceDiscovery] =
35-
## Two registrars in distinct buckets of the routing table rooted at `serviceId`.
36-
## Returned in the order `lookup()` would visit them (lower bucket index first).
37-
var a = setupServiceDiscoveryNode(discoConfig = conf)
38-
var b = setupServiceDiscoveryNode(discoConfig = conf)
39-
for _ in 0 ..< MaxRegistrarSetupAttempts:
40-
if bucketOf(a, serviceId) != bucketOf(b, serviceId):
41-
if bucketOf(a, serviceId) < bucketOf(b, serviceId):
42-
return (a, b)
43-
else:
44-
return (b, a)
45-
46-
b = setupServiceDiscoveryNode(discoConfig = conf)
47-
48-
raiseAssert "could not find registrars in distinct buckets"
49-
50-
proc setupRegistrarsInSameBucket(
51-
conf: ServiceDiscoveryConfig, serviceId: ServiceId, count: int
52-
): seq[ServiceDiscovery] =
53-
## Registrars that land in one service routing table bucket.
54-
doAssert count > 0, "count must be > 0"
55-
56-
var registrarsByBucket = initTable[int, seq[ServiceDiscovery]]()
57-
for _ in 0 ..< MaxRegistrarSetupAttempts:
58-
let registrar = setupServiceDiscoveryNode(discoConfig = conf)
59-
let bucket = bucketOf(registrar, serviceId)
60-
if bucket >= conf.bucketsCount:
61-
continue
62-
63-
if not registrarsByBucket.hasKey(bucket):
64-
registrarsByBucket[bucket] = @[]
65-
registrarsByBucket[bucket].add(registrar)
66-
67-
if registrarsByBucket[bucket].len == count:
68-
return registrarsByBucket[bucket]
69-
70-
raiseAssert "could not find enough registrars in one bucket"
71-
7221
suite "Service Discovery Component - Lookup Get Ads":
7322
teardown:
7423
checkTrackers()
@@ -155,27 +104,25 @@ suite "Service Discovery Component - Lookup Get Ads":
155104
)
156105
let discovererNode = setupServiceDiscoveryNode(discoConfig = conf)
157106

158-
let serviceName = "service"
107+
let (queriedFirst, queriedSecond, serviceName) =
108+
setupRegistrarsInDistinctBuckets(conf)
159109
let serviceId = serviceName.hashServiceId()
160-
let registrars = setupRegistrarsInDistinctBuckets(conf, serviceId)
161110

162-
startAndDeferStop(
163-
@[discovererNode, registrars.queriedFirst, registrars.queriedSecond]
164-
)
165-
await connect(discovererNode, registrars.queriedFirst)
166-
await connect(discovererNode, registrars.queriedSecond)
111+
startAndDeferStop(@[discovererNode, queriedFirst, queriedSecond])
112+
await connect(discovererNode, queriedFirst)
113+
await connect(discovererNode, queriedSecond)
167114

168115
# The first-queried registrar fills the result on its own.
169116
var firstBucketAds: seq[Advertisement]
170117
for _ in 0 ..< fLookup:
171118
firstBucketAds.add(makeAdvertisement(serviceName))
172-
registrars.queriedFirst.registrar.cache[serviceId] = firstBucketAds
119+
queriedFirst.registrar.cache[serviceId] = firstBucketAds
173120

174121
# The other should never be queried, so should not appear in the result.
175122
let otherKey = randomKey()
176123
let otherAd = makeAdvertisement(serviceName, privateKey = otherKey)
177124
let otherPeerId = PeerId.init(otherKey).get()
178-
registrars.queriedSecond.registrar.cache[serviceId] = @[otherAd]
125+
queriedSecond.registrar.cache[serviceId] = @[otherAd]
179126

180127
let found = await discovererNode.lookup(serviceId)
181128
check:
@@ -189,9 +136,8 @@ suite "Service Discovery Component - Lookup Get Ads":
189136
)
190137
let discovererNode = setupServiceDiscoveryNode(discoConfig = conf)
191138

192-
let serviceName = "service"
139+
let (registrars, serviceName) = setupRegistrarsInSameBucket(conf, kLookup + 2)
193140
let serviceId = serviceName.hashServiceId()
194-
let registrars = setupRegistrarsInSameBucket(conf, serviceId, kLookup + 2)
195141

196142
startAndDeferStop(@[discovererNode] & registrars)
197143
for registrar in registrars:
@@ -208,21 +154,19 @@ suite "Service Discovery Component - Lookup Get Ads":
208154
let conf = ServiceDiscoveryConfig.new(safetyParam = 0.0, fLookup = 1)
209155
let discovererNode = setupServiceDiscoveryNode(discoConfig = conf)
210156

211-
let serviceName = "service"
157+
let (queriedFirst, queriedSecond, serviceName) =
158+
setupRegistrarsInDistinctBuckets(conf)
212159
let serviceId = serviceName.hashServiceId()
213-
let registrars = setupRegistrarsInDistinctBuckets(conf, serviceId)
214160

215-
startAndDeferStop(
216-
@[discovererNode, registrars.queriedFirst, registrars.queriedSecond]
217-
)
218-
await connect(discovererNode, registrars.queriedFirst)
219-
await connect(discovererNode, registrars.queriedSecond)
161+
startAndDeferStop(@[discovererNode, queriedFirst, queriedSecond])
162+
await connect(discovererNode, queriedFirst)
163+
await connect(discovererNode, queriedSecond)
220164

221165
let firstBucketKey = randomKey()
222166
let firstBucketAd = makeAdvertisement(serviceName, privateKey = firstBucketKey)
223167
let secondBucketAd = makeAdvertisement(serviceName)
224-
registrars.queriedFirst.registrar.cache[serviceId] = @[firstBucketAd]
225-
registrars.queriedSecond.registrar.cache[serviceId] = @[secondBucketAd]
168+
queriedFirst.registrar.cache[serviceId] = @[firstBucketAd]
169+
queriedSecond.registrar.cache[serviceId] = @[secondBucketAd]
226170

227171
let found = await discovererNode.lookup(serviceId)
228172
check:

0 commit comments

Comments
 (0)