mirror of
https://github.com/logos-storage/logos-storage-nim-dht.git
synced 2026-08-24 22:19:14 +00:00
chore: update libp2p version (#116)
This commit is contained in:
@@ -34,14 +34,28 @@ func getField*(pb: ProtoBuffer, field: int,
|
||||
else:
|
||||
err(ProtoError.IncorrectBlob)
|
||||
|
||||
proc getField*(pb: ProtoBuffer, field: int,
|
||||
spr: var SignedPeerRecord): ProtoResult[bool] {.inline.} =
|
||||
## Read ``SignedPeerRecord`` from ProtoBuf's message and validate it
|
||||
var buffer: seq[byte]
|
||||
let res = ? pb.getField(field, buffer)
|
||||
if not(res):
|
||||
ok(false)
|
||||
else:
|
||||
let res2 = SignedPeerRecord.decode(buffer)
|
||||
if res2.isOk():
|
||||
spr = res2.get()
|
||||
ok(true)
|
||||
else:
|
||||
err(ProtoError.IncorrectBlob)
|
||||
|
||||
func write*[T: SignedPeerRecord | PeerRecord | Envelope](
|
||||
pb: var ProtoBuffer,
|
||||
field: int,
|
||||
env: T) {.raises: [Defect, ResultError[CryptoError]].} =
|
||||
env: T) {.raises: [Defect].} =
|
||||
|
||||
## Write Envelope value ``env`` to object ``pb`` using ProtoBuf's encoding.
|
||||
let encoded = env.encode().tryGet()
|
||||
write(pb, field, encoded)
|
||||
write(pb, field, env.encode())
|
||||
|
||||
proc getRepeatedField*(pb: ProtoBuffer, field: int,
|
||||
value: var seq[SignedPeerRecord]): ProtoResult[bool] {.
|
||||
|
||||
@@ -16,7 +16,6 @@
|
||||
import
|
||||
std/[hashes, net, options, sugar, tables],
|
||||
stew/endians2,
|
||||
bearssl/rand,
|
||||
chronicles,
|
||||
stew/[byteutils],
|
||||
stint,
|
||||
@@ -211,14 +210,14 @@ proc encodeStaticHeader*(flag: Flag, nonce: AESGCMNonce, authSize: int):
|
||||
# TODO: assert on authSize of > 2^16?
|
||||
result.add((uint16(authSize)).toBytesBE())
|
||||
|
||||
proc encodeMessagePacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
proc encodeMessagePacket*(rng: Rng, c: var Codec,
|
||||
toId: NodeId, toAddr: Address, message: openArray[byte]):
|
||||
(seq[byte], AESGCMNonce, bool) =
|
||||
var nonce: AESGCMNonce
|
||||
var haskey: bool
|
||||
hmacDrbgGenerate(rng, nonce) # Random AESGCM nonce
|
||||
rng.generate(nonce) # Random AESGCM nonce
|
||||
var iv: array[ivSize, byte]
|
||||
hmacDrbgGenerate(rng, iv) # Random IV
|
||||
rng.generate(iv) # Random IV
|
||||
|
||||
# static-header
|
||||
let authdata = c.localNode.id.toByteArrayBE()
|
||||
@@ -246,7 +245,7 @@ proc encodeMessagePacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
# case this must not look like a random packet.
|
||||
haskey = false
|
||||
var randomData: array[gcmTagSize + 4, byte]
|
||||
hmacDrbgGenerate(rng, randomData)
|
||||
rng.generate(randomData)
|
||||
messageEncrypted.add(randomData)
|
||||
dht_session_lru_cache_misses.inc()
|
||||
|
||||
@@ -259,11 +258,11 @@ proc encodeMessagePacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
|
||||
return (packet, nonce, haskey)
|
||||
|
||||
proc encodeWhoareyouPacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
proc encodeWhoareyouPacket*(rng: Rng, c: var Codec,
|
||||
toId: NodeId, toAddr: Address, requestNonce: AESGCMNonce, recordSeq: uint64,
|
||||
pubkey: Option[PublicKey]): seq[byte] =
|
||||
var idNonce: IdNonce
|
||||
hmacDrbgGenerate(rng, idNonce)
|
||||
rng.generate(idNonce)
|
||||
|
||||
# authdata
|
||||
var authdata: seq[byte]
|
||||
@@ -280,7 +279,7 @@ proc encodeWhoareyouPacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
header.add(authdata)
|
||||
|
||||
var iv: array[ivSize, byte]
|
||||
hmacDrbgGenerate(rng, iv) # Random IV
|
||||
rng.generate(iv) # Random IV
|
||||
|
||||
let maskedHeader = encryptHeader(toId, iv, header)
|
||||
|
||||
@@ -301,14 +300,14 @@ proc encodeWhoareyouPacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
|
||||
return packet
|
||||
|
||||
proc encodeHandshakePacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
proc encodeHandshakePacket*(rng: Rng, c: var Codec,
|
||||
toId: NodeId, toAddr: Address, message: openArray[byte],
|
||||
whoareyouData: WhoareyouData, pubkey: PublicKey): EncodeResult[seq[byte]] =
|
||||
var header: seq[byte]
|
||||
var nonce: AESGCMNonce
|
||||
hmacDrbgGenerate(rng, nonce)
|
||||
rng.generate(nonce)
|
||||
var iv: array[ivSize, byte]
|
||||
hmacDrbgGenerate(rng, iv) # Random IV
|
||||
rng.generate(iv) # Random IV
|
||||
|
||||
var authdata: seq[byte]
|
||||
var authdataHead: seq[byte]
|
||||
@@ -343,9 +342,7 @@ proc encodeHandshakePacket*(rng: var HmacDrbgContext, c: var Codec,
|
||||
|
||||
# Add SPR of sequence number is newer
|
||||
if whoareyouData.recordSeq < c.localNode.record.seqNum:
|
||||
let encoded = ? c.localNode.record.encode.mapErr((e: CryptoError) =>
|
||||
("Failed to encode local node's SignedPeerRecord: " & $e).cstring)
|
||||
authdata.add(encoded)
|
||||
authdata.add(c.localNode.record.encode)
|
||||
|
||||
let secrets = ? deriveKeys(
|
||||
c.localNode.id,
|
||||
|
||||
@@ -14,11 +14,12 @@
|
||||
|
||||
import
|
||||
std/[hashes, net],
|
||||
bearssl/rand,
|
||||
./spr,
|
||||
./node,
|
||||
../../../../dht/providers_messages
|
||||
|
||||
from libp2p/crypto/crypto import Rng, generate
|
||||
|
||||
export providers_messages
|
||||
|
||||
type
|
||||
@@ -131,7 +132,7 @@ template messageKind*(T: typedesc[SomeMessage]): MessageKind =
|
||||
proc hash*(reqId: RequestId): Hash =
|
||||
hash(reqId.id)
|
||||
|
||||
proc init*(T: type RequestId, rng: var HmacDrbgContext): T =
|
||||
proc init*(T: type RequestId, rng: Rng): T =
|
||||
var reqId = RequestId(id: newSeq[byte](8)) # RequestId must be <= 8 bytes
|
||||
hmacDrbgGenerate(rng, reqId.id)
|
||||
rng.generate(reqId.id)
|
||||
reqId
|
||||
|
||||
@@ -316,11 +316,7 @@ proc encodeMessage*[T: SomeMessage](p: T, reqId: RequestId, clientMode: bool = f
|
||||
result = newSeqOfCap[byte](64)
|
||||
result.add(messageKind(T).ord)
|
||||
|
||||
let encoded =
|
||||
try: p.encode()
|
||||
except ResultError[CryptoError] as e:
|
||||
error "Failed to encode protobuf message", typ = $T, msg = e.msg
|
||||
@[]
|
||||
let encoded = p.encode()
|
||||
var pb = initProtoBuffer()
|
||||
pb.write(1, reqId)
|
||||
pb.write(2, encoded)
|
||||
|
||||
@@ -9,7 +9,6 @@
|
||||
|
||||
import
|
||||
std/[hashes, net],
|
||||
bearssl/rand,
|
||||
chronicles,
|
||||
chronos,
|
||||
nimcrypto,
|
||||
@@ -17,6 +16,8 @@ import
|
||||
./crypto,
|
||||
./spr
|
||||
|
||||
from libp2p/crypto/crypto import Rng, generate
|
||||
|
||||
export stint
|
||||
|
||||
const
|
||||
@@ -132,9 +133,9 @@ func `==`*(a, b: Node): bool =
|
||||
func hash*(id: NodeId): Hash =
|
||||
hash(id.toByteArrayBE)
|
||||
|
||||
proc random*(T: type NodeId, rng: var HmacDrbgContext): T =
|
||||
proc random*(T: type NodeId, rng: Rng): T =
|
||||
var id: NodeId
|
||||
hmacDrbgGenerate(rng, addr id, csize_t(sizeof(id)))
|
||||
rng.generate(id)
|
||||
|
||||
id
|
||||
|
||||
|
||||
@@ -80,7 +80,6 @@ import
|
||||
pkg/[chronicles, chronicles/chronos_tools],
|
||||
pkg/chronos,
|
||||
pkg/stint,
|
||||
pkg/bearssl/rand,
|
||||
pkg/metrics,
|
||||
pkg/results
|
||||
|
||||
@@ -98,6 +97,8 @@ import "."/[
|
||||
|
||||
import nimcrypto except toHex
|
||||
|
||||
from libp2p/crypto/crypto import Rng
|
||||
|
||||
export options, results, node, spr, providers
|
||||
|
||||
declareCounter dht_message_requests_outgoing,
|
||||
@@ -177,7 +178,7 @@ type
|
||||
ipVote: IpVote
|
||||
enrAutoUpdate: bool
|
||||
talkProtocols*: Table[seq[byte], TalkProtocol] # TODO: Table is a bit of
|
||||
rng*: ref HmacDrbgContext
|
||||
rng*: Rng
|
||||
providers: ProvidersManager
|
||||
clientMode*: bool
|
||||
|
||||
@@ -483,7 +484,7 @@ proc sendRequest*[T: SomeMessage](d: Protocol, toNode: Node, m: T,
|
||||
|
||||
proc waitResponse*[T: SomeMessage](d: Protocol, node: Node, msg: T):
|
||||
Future[Option[Message]] =
|
||||
let reqId = RequestId.init(d.rng[])
|
||||
let reqId = RequestId.init(d.rng)
|
||||
result = d.waitMessage(node, reqId)
|
||||
sendRequest(d, node, msg, reqId)
|
||||
|
||||
@@ -500,7 +501,7 @@ proc waitMessage(d: Protocol, fromNode: Node, reqId: RequestId, timeout = Respon
|
||||
|
||||
proc waitNodeResponses*[T: SomeMessage](d: Protocol, node: Node, msg: T):
|
||||
Future[DiscResult[seq[SignedPeerRecord]]] =
|
||||
let reqId = RequestId.init(d.rng[])
|
||||
let reqId = RequestId.init(d.rng)
|
||||
result = d.waitNodes(node, reqId)
|
||||
sendRequest(d, node, msg, reqId)
|
||||
|
||||
@@ -742,7 +743,7 @@ proc addProvider*(
|
||||
res.add(d.localNode)
|
||||
for toNode in res:
|
||||
if toNode != d.localNode:
|
||||
let reqId = RequestId.init(d.rng[])
|
||||
let reqId = RequestId.init(d.rng)
|
||||
d.sendRequest(toNode, AddProviderMessage(cId: cId, prov: pr), reqId)
|
||||
else:
|
||||
asyncSpawn d.addProviderLocal(cId, pr)
|
||||
@@ -886,7 +887,7 @@ proc query*(d: Protocol, target: NodeId, k = BUCKET_SIZE): Future[seq[Node]]
|
||||
|
||||
proc queryRandom*(d: Protocol): Future[seq[Node]] =
|
||||
## Perform a query for a random target, return all nodes discovered.
|
||||
d.query(NodeId.random(d.rng[]))
|
||||
d.query(NodeId.random(d.rng))
|
||||
|
||||
proc queryRandom*(d: Protocol, enrField: (string, seq[byte])):
|
||||
Future[seq[Node]] {.async.} =
|
||||
@@ -978,7 +979,7 @@ proc revalidateLoop(d: Protocol) {.async.} =
|
||||
## message.
|
||||
try:
|
||||
while true:
|
||||
let revalidateTimeout = RevalidateMin + d.rng[].rand(RevalidateMax - RevalidateMin)
|
||||
let revalidateTimeout = RevalidateMin + d.rng.rand(RevalidateMax - RevalidateMin)
|
||||
await sleepAsync(milliseconds(revalidateTimeout))
|
||||
let n = d.routingTable.nodeToRevalidate()
|
||||
if not n.isNil:
|
||||
|
||||
@@ -10,7 +10,7 @@
|
||||
import std/sequtils
|
||||
|
||||
import pkg/chronicles
|
||||
import pkg/libp2p
|
||||
import pkg/libp2p/[peerid, routing_record]
|
||||
import pkg/questionable
|
||||
|
||||
import ../node
|
||||
|
||||
@@ -11,7 +11,7 @@ import std/sequtils
|
||||
import std/strutils
|
||||
|
||||
import pkg/chronos
|
||||
import pkg/libp2p
|
||||
import pkg/libp2p/[peerid, routing_record]
|
||||
import pkg/datastore
|
||||
import pkg/questionable
|
||||
import pkg/questionable/results
|
||||
|
||||
@@ -13,7 +13,7 @@ from std/times import now, utc, toTime, toUnix
|
||||
|
||||
import pkg/stew/endians2
|
||||
import pkg/chronos
|
||||
import pkg/libp2p
|
||||
import pkg/libp2p/[peerid, routing_record]
|
||||
import pkg/datastore
|
||||
import pkg/chronicles
|
||||
import pkg/questionable
|
||||
|
||||
@@ -12,9 +12,8 @@ from std/times import now, utc, toTime, toUnix
|
||||
import pkg/stew/endians2
|
||||
import pkg/datastore
|
||||
import pkg/chronos
|
||||
import pkg/libp2p
|
||||
import pkg/libp2p/[peerid, routing_record]
|
||||
import pkg/chronicles
|
||||
import pkg/stew/byteutils
|
||||
import pkg/questionable
|
||||
import pkg/questionable/results
|
||||
|
||||
@@ -88,10 +87,7 @@ proc add*(
|
||||
trace "Provider with same seqNo already exist", seqNo = $provider.data.seqNo
|
||||
@[]
|
||||
else:
|
||||
without bytes =? provider.envelope.encode:
|
||||
trace "Enable to encode provider"
|
||||
return failure "Unable to encode provider"
|
||||
bytes
|
||||
provider.envelope.encode
|
||||
|
||||
if bytes.len > 0:
|
||||
trace "Adding or updating provider record", id, peerId
|
||||
|
||||
@@ -1,22 +1,22 @@
|
||||
import bearssl/rand
|
||||
from libp2p/crypto/crypto import Rng, generate
|
||||
|
||||
## Random helpers: similar as in stdlib, but with HmacDrbgContext rng
|
||||
## Random helpers: similar as in stdlib, but with libp2p Rng
|
||||
# TODO: Move these somewhere else?
|
||||
const randMax = 18_446_744_073_709_551_615'u64
|
||||
|
||||
proc rand*(rng: var HmacDrbgContext, max: Natural): int =
|
||||
proc rand*(rng: Rng, max: Natural): int =
|
||||
if max == 0: return 0
|
||||
|
||||
var x: uint64
|
||||
while true:
|
||||
hmacDrbgGenerate(rng, addr x, csize_t(sizeof(x)))
|
||||
rng.generate(x)
|
||||
if x < randMax - (randMax mod (uint64(max) + 1'u64)): # against modulo bias
|
||||
return int(x mod (uint64(max) + 1'u64))
|
||||
|
||||
proc sample*[T](rng: var HmacDrbgContext, a: openArray[T]): T =
|
||||
proc sample*[T](rng: Rng, a: openArray[T]): T =
|
||||
result = a[rng.rand(a.high)]
|
||||
|
||||
proc shuffle*[T](rng: var HmacDrbgContext, a: var openArray[T]) =
|
||||
proc shuffle*[T](rng: Rng, a: var openArray[T]) =
|
||||
for i in countdown(a.high, 1):
|
||||
let j = rng.rand(i)
|
||||
swap(a[i], a[j])
|
||||
|
||||
@@ -9,9 +9,11 @@
|
||||
|
||||
import
|
||||
std/[algorithm, net, times, sequtils, bitops, sets, options, tables],
|
||||
stint, chronicles, metrics, bearssl/rand, chronos,
|
||||
stint, chronicles, metrics, chronos,
|
||||
"."/[node, random2, spr]
|
||||
|
||||
from libp2p/crypto/crypto import Rng
|
||||
|
||||
export options
|
||||
|
||||
declarePublicGauge dht_routing_table_nodes,
|
||||
@@ -51,7 +53,7 @@ type
|
||||
ipLimits: IpLimits ## IP limits for total routing table: all buckets and
|
||||
## replacement caches.
|
||||
distanceCalculator: DistanceCalculator
|
||||
rng: ref HmacDrbgContext
|
||||
rng: Rng
|
||||
|
||||
KBucket = ref object
|
||||
istart, iend: NodeId ## Range of NodeIds this KBucket covers. This is not a
|
||||
@@ -289,7 +291,7 @@ proc getDepth*(b: KBucket) : int =
|
||||
computeSharedPrefixBits(@[b.istart, b.iend])
|
||||
|
||||
proc init*(T: type RoutingTable, localNode: Node, bitsPerHop = DefaultBitsPerHop,
|
||||
ipLimits = DefaultTableIpLimits, rng: ref HmacDrbgContext,
|
||||
ipLimits = DefaultTableIpLimits, rng: Rng,
|
||||
distanceCalculator = XorDistanceCalculator): T =
|
||||
## Initialize the routing table for provided `Node` and bitsPerHop value.
|
||||
## `bitsPerHop` is default set to 5 as recommended by original Kademlia paper.
|
||||
@@ -553,7 +555,7 @@ proc nodeToRevalidate*(r: RoutingTable): Node =
|
||||
## Return a node to revalidate. The least recently seen node from a random
|
||||
## bucket is selected.
|
||||
var buckets = r.buckets
|
||||
r.rng[].shuffle(buckets)
|
||||
random2.shuffle(r.rng, buckets)
|
||||
# TODO: Should we prioritize less-recently-updated buckets instead? Could
|
||||
# store a `now` Moment at setJustSeen or at revalidate per bucket.
|
||||
for b in buckets:
|
||||
@@ -586,9 +588,9 @@ proc randomNodes*(r: RoutingTable, maxAmount: int,
|
||||
# already.
|
||||
# We check against the number of nodes to avoid an infinite loop in case of a filter.
|
||||
while len(result) < maxAmount and len(seen) < sz:
|
||||
let bucket = r.rng[].sample(r.buckets)
|
||||
let bucket = r.rng.sample(r.buckets)
|
||||
if bucket.nodes.len != 0:
|
||||
let node = r.rng[].sample(bucket.nodes)
|
||||
let node = r.rng.sample(bucket.nodes)
|
||||
if node notin seen:
|
||||
seen.incl(node)
|
||||
if pred.isNil() or node.pred:
|
||||
|
||||
@@ -14,7 +14,8 @@ import
|
||||
libp2p/crypto/crypto,
|
||||
libp2p/crypto/secp,
|
||||
libp2p/routing_record,
|
||||
libp2p/multicodec
|
||||
libp2p/multicodec,
|
||||
pkg/protobuf_serialization
|
||||
|
||||
export routing_record
|
||||
|
||||
@@ -98,7 +99,7 @@ proc update*(
|
||||
return err "No existing address in SignedPeerRecord with no port provided"
|
||||
|
||||
let ipAddr = ip.get
|
||||
|
||||
|
||||
if tcpPort.isSome:
|
||||
transProto = IpTransportProtocol.tcpProtocol
|
||||
transProtoPort = tcpPort.get
|
||||
@@ -215,10 +216,7 @@ template fromURI*(r: var SignedPeerRecord, url: SprUri): bool =
|
||||
fromURI(r, string(url))
|
||||
|
||||
proc toBase64*(r: SignedPeerRecord): string =
|
||||
let encoded = r.encode
|
||||
if encoded.isErr:
|
||||
error "Failed to encode SignedPeerRecord", error = encoded.error
|
||||
result = Base64Url.encode(encoded.get(@[]))
|
||||
Base64Url.encode(r.encode)
|
||||
|
||||
proc toURI*(r: SignedPeerRecord): string = "spr:" & r.toBase64
|
||||
|
||||
|
||||
@@ -7,7 +7,6 @@
|
||||
# Everything below the handling of ordinary messages
|
||||
import
|
||||
std/[net, tables, options, sets],
|
||||
bearssl/rand,
|
||||
chronos,
|
||||
chronicles,
|
||||
metrics,
|
||||
@@ -41,7 +40,7 @@ type
|
||||
keyexchangeInProgress: HashSet[NodeId]
|
||||
pendingRequestsByNode: Table[NodeId, seq[seq[byte]]]
|
||||
codec*: Codec
|
||||
rng: ref HmacDrbgContext
|
||||
rng: Rng
|
||||
|
||||
PendingRequest = object
|
||||
node: Node
|
||||
@@ -76,7 +75,7 @@ proc send(t: Transport, n: Node, data: seq[byte]) =
|
||||
t.sendToA(n.address.get(), data)
|
||||
|
||||
proc sendMessage*(t: Transport, toId: NodeId, toAddr: Address, message: seq[byte]) =
|
||||
let (data, _, _) = encodeMessagePacket(t.rng[], t.codec, toId, toAddr,
|
||||
let (data, _, _) = encodeMessagePacket(t.rng, t.codec, toId, toAddr,
|
||||
message)
|
||||
t.sendToA(toAddr, data)
|
||||
|
||||
@@ -94,7 +93,7 @@ proc registerRequest(t: Transport, n: Node, message: seq[byte],
|
||||
proc sendMessage*(t: Transport, toNode: Node, message: seq[byte]) =
|
||||
doAssert(toNode.address.isSome())
|
||||
let address = toNode.address.get()
|
||||
let (data, nonce, haskey) = encodeMessagePacket(t.rng[], t.codec,
|
||||
let (data, nonce, haskey) = encodeMessagePacket(t.rng, t.codec,
|
||||
toNode.id, address, message)
|
||||
|
||||
if haskey:
|
||||
@@ -129,7 +128,7 @@ proc sendWhoareyou(t: Transport, toId: NodeId, a: Address,
|
||||
pubkey = if node.isSome(): some(node.get().pubkey)
|
||||
else: none(PublicKey)
|
||||
|
||||
let data = encodeWhoareyouPacket(t.rng[], t.codec, toId, a, requestNonce,
|
||||
let data = encodeWhoareyouPacket(t.rng, t.codec, toId, a, requestNonce,
|
||||
recordSeq, pubkey)
|
||||
sleepAsync(handshakeTimeout).addCallback() do(data: pointer):
|
||||
# handshake key is popped in decodeHandshakePacket. if not yet popped by timeout:
|
||||
@@ -151,7 +150,7 @@ proc sendPending(t:Transport, toNode: Node):
|
||||
for message in t.pendingRequestsByNode[toNode.id]:
|
||||
trace "Sending pending packet", myport = t.bindAddress.port, dstId = toNode.id
|
||||
let address = toNode.address.get()
|
||||
let (data, nonce, haskey) = encodeMessagePacket(t.rng[], t.codec, toNode.id, address, message)
|
||||
let (data, nonce, haskey) = encodeMessagePacket(t.rng, t.codec, toNode.id, address, message)
|
||||
t.registerRequest(toNode, message, nonce)
|
||||
t.send(toNode, data)
|
||||
t.pendingRequestsByNode.del(toNode.id)
|
||||
@@ -197,7 +196,7 @@ proc receive*(t: Transport, a: Address, packet: openArray[byte]) =
|
||||
doAssert(toNode.address.isSome())
|
||||
let address = toNode.address.get()
|
||||
let data = encodeHandshakePacket(
|
||||
t.rng[],
|
||||
t.rng,
|
||||
t.codec,
|
||||
toNode.id,
|
||||
address,
|
||||
|
||||
Reference in New Issue
Block a user