2022-07-01 20:19:57 +02:00
|
|
|
# Nim-LibP2P
|
2023-01-20 15:47:40 +01:00
|
|
|
# Copyright (c) 2023 Status Research & Development GmbH
|
2022-07-01 20:19:57 +02:00
|
|
|
# Licensed under either of
|
|
|
|
# * Apache License, version 2.0, ([LICENSE-APACHE](LICENSE-APACHE))
|
|
|
|
# * MIT license ([LICENSE-MIT](LICENSE-MIT))
|
|
|
|
# at your option.
|
|
|
|
# This file may not be copied, modified, or distributed except according to
|
|
|
|
# those terms.
|
2019-09-09 20:15:52 -06:00
|
|
|
|
2023-06-07 13:12:49 +02:00
|
|
|
{.push raises: [].}
|
2021-05-21 10:27:01 -06:00
|
|
|
|
2023-04-03 10:56:20 +02:00
|
|
|
import std/[sequtils, strutils, tables, hashes, options, sets, deques]
|
2022-09-22 21:55:59 +02:00
|
|
|
import stew/results
|
2020-06-11 12:09:34 -06:00
|
|
|
import chronos, chronicles, nimcrypto/sha2, metrics
|
2023-09-22 16:45:08 +02:00
|
|
|
import chronos/ratelimit
|
2019-12-10 14:50:35 -06:00
|
|
|
import rpc/[messages, message, protobuf],
|
2020-07-01 15:25:09 +09:00
|
|
|
../../peerid,
|
2019-09-24 10:16:39 -06:00
|
|
|
../../peerinfo,
|
2020-06-19 11:29:43 -06:00
|
|
|
../../stream/connection,
|
2019-09-11 23:46:08 -06:00
|
|
|
../../crypto/crypto,
|
2020-03-23 15:03:36 +09:00
|
|
|
../../protobuf/minprotobuf,
|
|
|
|
../../utility
|
2019-09-09 20:15:52 -06:00
|
|
|
|
2023-07-28 10:58:05 +02:00
|
|
|
export peerid, connection, deques
|
2020-09-24 18:43:20 +02:00
|
|
|
|
2019-09-09 20:15:52 -06:00
|
|
|
logScope:
|
2020-12-01 11:34:27 -06:00
|
|
|
topics = "libp2p pubsubpeer"
|
2019-09-09 20:15:52 -06:00
|
|
|
|
2020-08-05 01:27:59 +02:00
|
|
|
when defined(libp2p_expensive_metrics):
|
|
|
|
declareCounter(libp2p_pubsub_sent_messages, "number of messages sent", labels = ["id", "topic"])
|
|
|
|
declareCounter(libp2p_pubsub_skipped_received_messages, "number of received skipped messages", labels = ["id"])
|
|
|
|
declareCounter(libp2p_pubsub_skipped_sent_messages, "number of sent skipped messages", labels = ["id"])
|
2020-06-11 20:20:58 -06:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
when defined(pubsubpeer_queue_metrics):
|
|
|
|
declareGauge(libp2p_gossipsub_priority_queue_size, "the number of messages in the priority queue", labels = ["id"])
|
|
|
|
declareGauge(libp2p_gossipsub_non_priority_queue_size, "the number of messages in the non-priority queue", labels = ["id"])
|
|
|
|
|
2024-03-25 22:00:11 +01:00
|
|
|
declareCounter(libp2p_pubsub_disconnects_over_non_priority_queue_limit, "number of peers disconnected due to over non-prio queue capacity")
|
|
|
|
|
|
|
|
const
|
|
|
|
DefaultMaxNumElementsInNonPriorityQueue* = 1024
|
|
|
|
|
2019-09-09 20:15:52 -06:00
|
|
|
type
|
2023-09-22 16:45:08 +02:00
|
|
|
PeerRateLimitError* = object of CatchableError
|
|
|
|
|
2020-04-30 22:22:31 +09:00
|
|
|
PubSubObserver* = ref object
|
2023-06-07 13:12:49 +02:00
|
|
|
onRecv*: proc(peer: PubSubPeer; msgs: var RPCMsg) {.gcsafe, raises: [].}
|
|
|
|
onSend*: proc(peer: PubSubPeer; msgs: var RPCMsg) {.gcsafe, raises: [].}
|
2019-09-09 20:15:52 -06:00
|
|
|
|
2020-09-22 09:05:53 +02:00
|
|
|
PubSubPeerEventKind* {.pure.} = enum
|
2024-03-25 22:00:11 +01:00
|
|
|
StreamOpened
|
|
|
|
StreamClosed
|
|
|
|
DisconnectionRequested # tells gossipsub that the transport connection to the peer should be closed
|
2020-09-22 09:05:53 +02:00
|
|
|
|
2022-07-27 17:14:05 +00:00
|
|
|
PubSubPeerEvent* = object
|
2020-09-22 09:05:53 +02:00
|
|
|
kind*: PubSubPeerEventKind
|
|
|
|
|
2023-06-07 13:12:49 +02:00
|
|
|
GetConn* = proc(): Future[Connection] {.gcsafe, raises: [].}
|
|
|
|
DropConn* = proc(peer: PubSubPeer) {.gcsafe, raises: [].} # have to pass peer as it's unknown during init
|
|
|
|
OnEvent* = proc(peer: PubSubPeer, event: PubSubPeerEvent) {.gcsafe, raises: [].}
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
RpcMessageQueue* = ref object
|
|
|
|
# Tracks async tasks for sending high-priority peer-published messages.
|
|
|
|
sendPriorityQueue: Deque[Future[void]]
|
|
|
|
# Queue for lower-priority messages, like "IWANT" replies and relay messages.
|
|
|
|
nonPriorityQueue: AsyncQueue[seq[byte]]
|
|
|
|
# Task for processing non-priority message queue.
|
|
|
|
sendNonPriorityTask: Future[void]
|
|
|
|
|
2020-04-30 22:22:31 +09:00
|
|
|
PubSubPeer* = ref object of RootObj
|
2020-09-01 09:33:03 +02:00
|
|
|
getConn*: GetConn # callback to establish a new send connection
|
2020-09-22 09:05:53 +02:00
|
|
|
onEvent*: OnEvent # Connectivity updates for peer
|
2020-08-11 18:05:49 -06:00
|
|
|
codec*: string # the protocol that this peer joined from
|
2020-09-22 09:05:53 +02:00
|
|
|
sendConn*: Connection # cached send connection
|
2023-01-10 13:33:14 +01:00
|
|
|
connectedFut: Future[void]
|
2021-02-27 23:49:56 +09:00
|
|
|
address*: Option[MultiAddress]
|
2021-12-16 11:05:20 +01:00
|
|
|
peerId*: PeerId
|
2020-04-30 22:22:31 +09:00
|
|
|
handler*: RPCHandler
|
|
|
|
observers*: ref seq[PubSubObserver] # ref as in smart_ptr
|
|
|
|
|
2020-09-21 18:16:29 +09:00
|
|
|
score*: float64
|
2023-04-03 10:56:20 +02:00
|
|
|
sentIHaves*: Deque[HashSet[MessageId]]
|
2024-05-02 12:18:55 +02:00
|
|
|
heDontWants*: Deque[HashSet[SaltedId]]
|
|
|
|
## IDONTWANT contains unvalidated message id:s which may be long and/or
|
|
|
|
## expensive to look up, so we apply the same salting to them as during
|
|
|
|
## unvalidated message processing
|
2020-09-21 18:16:29 +09:00
|
|
|
iHaveBudget*: int
|
2023-06-21 10:40:10 +02:00
|
|
|
pingBudget*: int
|
2021-10-25 12:58:38 +02:00
|
|
|
maxMessageSize: int
|
2020-09-21 18:16:29 +09:00
|
|
|
appScore*: float64 # application specific score
|
|
|
|
behaviourPenalty*: float64 # the eventual penalty score
|
2023-09-22 16:45:08 +02:00
|
|
|
overheadRateLimitOpt*: Opt[TokenBucket]
|
2021-01-15 13:48:03 +09:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
rpcmessagequeue: RpcMessageQueue
|
2024-03-25 22:00:11 +01:00
|
|
|
maxNumElementsInNonPriorityQueue*: int # The max number of elements allowed in the non-priority queue.
|
|
|
|
disconnected: bool
|
2024-03-05 16:05:21 +01:00
|
|
|
|
2023-09-22 16:45:08 +02:00
|
|
|
RPCHandler* = proc(peer: PubSubPeer, data: seq[byte]): Future[void]
|
2023-06-07 13:12:49 +02:00
|
|
|
{.gcsafe, raises: [].}
|
2019-09-09 20:15:52 -06:00
|
|
|
|
2023-04-03 11:05:01 +02:00
|
|
|
when defined(libp2p_agents_metrics):
|
|
|
|
func shortAgent*(p: PubSubPeer): string =
|
|
|
|
if p.sendConn.isNil or p.sendConn.getWrapped().isNil:
|
|
|
|
"unknown"
|
|
|
|
else:
|
|
|
|
#TODO the sendConn is setup before identify,
|
|
|
|
#so we have to read the parents short agent..
|
|
|
|
p.sendConn.getWrapped().shortAgent
|
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
proc getAgent*(peer: PubSubPeer): string =
|
|
|
|
return
|
|
|
|
when defined(libp2p_agents_metrics):
|
|
|
|
if peer.shortAgent.len > 0:
|
|
|
|
peer.shortAgent
|
|
|
|
else:
|
|
|
|
"unknown"
|
|
|
|
else:
|
|
|
|
"unknown"
|
|
|
|
|
2020-09-22 09:05:53 +02:00
|
|
|
func hash*(p: PubSubPeer): Hash =
|
2021-02-26 19:19:15 +09:00
|
|
|
p.peerId.hash
|
|
|
|
|
|
|
|
func `==`*(a, b: PubSubPeer): bool =
|
|
|
|
a.peerId == b.peerId
|
2020-07-13 22:32:38 +09:00
|
|
|
|
2020-09-06 10:31:47 +02:00
|
|
|
func shortLog*(p: PubSubPeer): string =
|
|
|
|
if p.isNil: "PubSubPeer(nil)"
|
|
|
|
else: shortLog(p.peerId)
|
|
|
|
chronicles.formatIt(PubSubPeer): shortLog(it)
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2020-06-29 09:15:31 -06:00
|
|
|
proc connected*(p: PubSubPeer): bool =
|
2020-08-11 18:05:49 -06:00
|
|
|
not p.sendConn.isNil and not
|
|
|
|
(p.sendConn.closed or p.sendConn.atEof)
|
2020-07-07 18:33:05 -06:00
|
|
|
|
2021-04-18 10:08:33 +02:00
|
|
|
proc hasObservers*(p: PubSubPeer): bool =
|
2020-09-24 18:43:20 +02:00
|
|
|
p.observers != nil and anyIt(p.observers[], it != nil)
|
|
|
|
|
2021-02-06 09:13:04 +09:00
|
|
|
func outbound*(p: PubSubPeer): bool =
|
2021-03-03 08:23:40 +09:00
|
|
|
# gossipsub 1.1 spec requires us to know if the transport is outgoing
|
|
|
|
# in order to give priotity to connections we make
|
|
|
|
# https://github.com/libp2p/specs/blob/master/pubsub/gossipsub/gossipsub-v1.1.md#outbound-mesh-quotas
|
|
|
|
# This behaviour is presrcibed to counter sybil attacks and ensures that a coordinated inbound attack can never fully take over the mesh
|
|
|
|
if not p.sendConn.isNil and p.sendConn.transportDir == Direction.Out:
|
2021-02-06 09:13:04 +09:00
|
|
|
true
|
|
|
|
else:
|
|
|
|
false
|
|
|
|
|
2023-09-22 16:45:08 +02:00
|
|
|
proc recvObservers*(p: PubSubPeer, msg: var RPCMsg) =
|
2020-05-29 10:46:27 -06:00
|
|
|
# trigger hooks
|
|
|
|
if not(isNil(p.observers)) and p.observers[].len > 0:
|
|
|
|
for obs in p.observers[]:
|
2020-06-24 09:08:44 -06:00
|
|
|
if not(isNil(obs)): # TODO: should never be nil, but...
|
|
|
|
obs.onRecv(p, msg)
|
2020-05-29 10:46:27 -06:00
|
|
|
|
|
|
|
proc sendObservers(p: PubSubPeer, msg: var RPCMsg) =
|
|
|
|
# trigger hooks
|
|
|
|
if not(isNil(p.observers)) and p.observers[].len > 0:
|
|
|
|
for obs in p.observers[]:
|
2020-06-24 09:08:44 -06:00
|
|
|
if not(isNil(obs)): # TODO: should never be nil, but...
|
|
|
|
obs.onSend(p, msg)
|
2020-05-29 10:46:27 -06:00
|
|
|
|
2019-12-16 23:24:03 -06:00
|
|
|
proc handle*(p: PubSubPeer, conn: Connection) {.async.} =
|
2020-09-06 10:31:47 +02:00
|
|
|
debug "starting pubsub read loop",
|
|
|
|
conn, peer = p, closed = conn.closed
|
2019-09-09 20:15:52 -06:00
|
|
|
try:
|
2020-05-25 08:33:24 -06:00
|
|
|
try:
|
2020-08-02 23:20:11 -06:00
|
|
|
while not conn.atEof:
|
2020-09-06 10:31:47 +02:00
|
|
|
trace "waiting for data", conn, peer = p, closed = conn.closed
|
|
|
|
|
2021-10-25 12:58:38 +02:00
|
|
|
var data = await conn.readLp(p.maxMessageSize)
|
2020-09-06 10:31:47 +02:00
|
|
|
trace "read data from peer",
|
|
|
|
conn, peer = p, closed = conn.closed,
|
|
|
|
data = data.shortLog
|
2020-05-25 08:33:24 -06:00
|
|
|
|
2023-09-22 16:45:08 +02:00
|
|
|
await p.handler(p, data)
|
|
|
|
data = newSeq[byte]() # Release memory
|
|
|
|
except PeerRateLimitError as exc:
|
|
|
|
debug "Peer rate limit exceeded, exiting read while", conn, peer = p, error = exc.msg
|
|
|
|
except CatchableError as exc:
|
|
|
|
debug "Exception occurred in PubSubPeer.handle",
|
|
|
|
conn, peer = p, closed = conn.closed, exc = exc.msg
|
2020-05-25 08:33:24 -06:00
|
|
|
finally:
|
|
|
|
await conn.close()
|
2020-09-04 19:30:45 +03:00
|
|
|
except CancelledError:
|
|
|
|
# This is top-level procedure which will work as separate task, so it
|
2020-11-23 15:02:23 -06:00
|
|
|
# do not need to propagate CancelledError.
|
2020-09-04 19:30:45 +03:00
|
|
|
trace "Unexpected cancellation in PubSubPeer.handle"
|
2019-12-05 20:16:18 -06:00
|
|
|
except CatchableError as exc:
|
2020-09-06 10:31:47 +02:00
|
|
|
trace "Exception occurred in PubSubPeer.handle",
|
|
|
|
conn, peer = p, closed = conn.closed, exc = exc.msg
|
2020-09-01 09:33:03 +02:00
|
|
|
finally:
|
2020-09-06 10:31:47 +02:00
|
|
|
debug "exiting pubsub read loop",
|
|
|
|
conn, peer = p, closed = conn.closed
|
2019-09-09 20:15:52 -06:00
|
|
|
|
2024-03-25 22:00:11 +01:00
|
|
|
proc closeSendConn(p: PubSubPeer, event: PubSubPeerEventKind) {.async.} =
|
|
|
|
if p.sendConn != nil:
|
|
|
|
trace "Removing send connection", p, conn = p.sendConn
|
|
|
|
await p.sendConn.close()
|
|
|
|
p.sendConn = nil
|
|
|
|
|
|
|
|
if not p.connectedFut.finished:
|
|
|
|
p.connectedFut.complete()
|
|
|
|
|
|
|
|
try:
|
|
|
|
if p.onEvent != nil:
|
|
|
|
p.onEvent(p, PubSubPeerEvent(kind: event))
|
|
|
|
except CancelledError as exc:
|
|
|
|
raise exc
|
|
|
|
except CatchableError as exc:
|
|
|
|
debug "Errors during diconnection events", error = exc.msg
|
|
|
|
# don't cleanup p.address else we leak some gossip stat table
|
|
|
|
|
2020-09-22 09:05:53 +02:00
|
|
|
proc connectOnce(p: PubSubPeer): Future[void] {.async.} =
|
2020-08-17 22:39:39 +03:00
|
|
|
try:
|
2023-01-10 13:33:14 +01:00
|
|
|
if p.connectedFut.finished:
|
|
|
|
p.connectedFut = newFuture[void]()
|
2023-06-14 17:23:39 +02:00
|
|
|
let newConn = await p.getConn().wait(5.seconds)
|
2020-08-17 22:39:39 +03:00
|
|
|
if newConn.isNil:
|
2021-06-02 07:39:10 -06:00
|
|
|
raise (ref LPError)(msg: "Cannot establish send connection")
|
2020-08-17 22:39:39 +03:00
|
|
|
|
2020-09-22 09:05:53 +02:00
|
|
|
# When the send channel goes up, subscriptions need to be sent to the
|
|
|
|
# remote peer - if we had multiple channels up and one goes down, all
|
|
|
|
# stop working so we make an effort to only keep a single channel alive
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2020-09-22 09:05:53 +02:00
|
|
|
trace "Get new send connection", p, newConn
|
2023-03-08 12:30:19 +01:00
|
|
|
|
|
|
|
# Careful to race conditions here.
|
|
|
|
# Topic subscription relies on either connectedFut
|
|
|
|
# to be completed, or onEvent to be called later
|
2023-01-10 13:33:14 +01:00
|
|
|
p.connectedFut.complete()
|
2020-08-17 22:39:39 +03:00
|
|
|
p.sendConn = newConn
|
2022-09-22 21:55:59 +02:00
|
|
|
p.address = if p.sendConn.observedAddr.isSome: some(p.sendConn.observedAddr.get) else: none(MultiAddress)
|
2020-09-22 09:05:53 +02:00
|
|
|
|
|
|
|
if p.onEvent != nil:
|
2024-03-25 22:00:11 +01:00
|
|
|
p.onEvent(p, PubSubPeerEvent(kind: PubSubPeerEventKind.StreamOpened))
|
2020-09-22 09:05:53 +02:00
|
|
|
|
|
|
|
await handle(p, newConn)
|
2020-08-17 22:39:39 +03:00
|
|
|
finally:
|
2024-03-25 22:00:11 +01:00
|
|
|
await p.closeSendConn(PubSubPeerEventKind.StreamClosed)
|
2020-09-22 09:05:53 +02:00
|
|
|
|
|
|
|
proc connectImpl(p: PubSubPeer) {.async.} =
|
2020-09-01 09:33:03 +02:00
|
|
|
try:
|
2020-09-22 09:05:53 +02:00
|
|
|
# Keep trying to establish a connection while it's possible to do so - the
|
|
|
|
# send connection might get disconnected due to a timeout or an unrelated
|
|
|
|
# issue so we try to get a new on
|
|
|
|
while true:
|
2024-03-25 22:00:11 +01:00
|
|
|
if p.disconnected:
|
|
|
|
if not p.connectedFut.finished:
|
|
|
|
p.connectedFut.complete()
|
|
|
|
return
|
2020-09-22 09:05:53 +02:00
|
|
|
await connectOnce(p)
|
2021-03-09 13:22:52 +01:00
|
|
|
except CatchableError as exc: # never cancelled
|
2020-09-22 09:05:53 +02:00
|
|
|
debug "Could not establish send connection", msg = exc.msg
|
2020-09-01 09:33:03 +02:00
|
|
|
|
|
|
|
proc connect*(p: PubSubPeer) =
|
2023-01-10 13:33:14 +01:00
|
|
|
if p.connected:
|
|
|
|
return
|
2020-07-16 12:06:57 +02:00
|
|
|
|
2023-01-10 13:33:14 +01:00
|
|
|
asyncSpawn connectImpl(p)
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2023-03-08 12:30:19 +01:00
|
|
|
proc hasSendConn*(p: PubSubPeer): bool =
|
|
|
|
p.sendConn != nil
|
|
|
|
|
2020-11-23 15:02:23 -06:00
|
|
|
template sendMetrics(msg: RPCMsg): untyped =
|
|
|
|
when defined(libp2p_expensive_metrics):
|
|
|
|
for x in msg.messages:
|
2024-03-25 12:06:34 +01:00
|
|
|
# metrics
|
|
|
|
libp2p_pubsub_sent_messages.inc(labelValues = [$p.peerId, x.topic])
|
2020-11-23 15:02:23 -06:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
proc clearSendPriorityQueue(p: PubSubPeer) =
|
|
|
|
if p.rpcmessagequeue.sendPriorityQueue.len == 0:
|
|
|
|
return # fast path
|
2024-02-19 13:47:37 +01:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
while p.rpcmessagequeue.sendPriorityQueue.len > 0 and
|
|
|
|
p.rpcmessagequeue.sendPriorityQueue[0].finished:
|
|
|
|
discard p.rpcmessagequeue.sendPriorityQueue.popFirst()
|
2024-02-19 13:47:37 +01:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
while p.rpcmessagequeue.sendPriorityQueue.len > 0 and
|
|
|
|
p.rpcmessagequeue.sendPriorityQueue[^1].finished:
|
|
|
|
discard p.rpcmessagequeue.sendPriorityQueue.popLast()
|
|
|
|
|
|
|
|
when defined(pubsubpeer_queue_metrics):
|
|
|
|
libp2p_gossipsub_priority_queue_size.set(
|
|
|
|
value = p.rpcmessagequeue.sendPriorityQueue.len.int64,
|
|
|
|
labelValues = [$p.peerId])
|
|
|
|
|
|
|
|
proc sendMsgContinue(conn: Connection, msgFut: Future[void]) {.async.} =
|
|
|
|
# Continuation for a pending `sendMsg` future from below
|
|
|
|
try:
|
|
|
|
await msgFut
|
|
|
|
trace "sent pubsub message to remote", conn
|
|
|
|
except CatchableError as exc: # never cancelled
|
|
|
|
# Because we detach the send call from the currently executing task using
|
|
|
|
# asyncSpawn, no exceptions may leak out of it
|
|
|
|
trace "Unable to send to remote", conn, msg = exc.msg
|
|
|
|
# Next time sendConn is used, it will be have its close flag set and thus
|
|
|
|
# will be recycled
|
|
|
|
|
|
|
|
await conn.close() # This will clean up the send connection
|
2021-10-25 12:58:38 +02:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
proc sendMsgSlow(p: PubSubPeer, msg: seq[byte]) {.async.} =
|
|
|
|
# Slow path of `sendMsg` where msg is held in memory while send connection is
|
|
|
|
# being set up
|
2023-01-10 13:33:14 +01:00
|
|
|
if p.sendConn == nil:
|
2023-06-14 17:23:39 +02:00
|
|
|
# Wait for a send conn to be setup. `connectOnce` will
|
|
|
|
# complete this even if the sendConn setup failed
|
2024-05-07 15:44:14 +02:00
|
|
|
discard await race(p.connectedFut)
|
2023-01-10 13:33:14 +01:00
|
|
|
|
|
|
|
var conn = p.sendConn
|
2020-09-24 18:43:20 +02:00
|
|
|
if conn == nil or conn.closed():
|
2023-06-14 17:23:39 +02:00
|
|
|
debug "No send connection", p, msg = shortLog(msg)
|
2020-09-24 18:43:20 +02:00
|
|
|
return
|
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
trace "sending encoded msg to peer", conn, encoded = shortLog(msg)
|
|
|
|
await sendMsgContinue(conn, conn.writeLp(msg))
|
2023-01-10 13:33:14 +01:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
proc sendMsg(p: PubSubPeer, msg: seq[byte]): Future[void] =
|
|
|
|
if p.sendConn != nil and not p.sendConn.closed():
|
|
|
|
# Fast path that avoids copying msg (which happens for {.async.})
|
|
|
|
let conn = p.sendConn
|
2023-01-10 13:33:14 +01:00
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
trace "sending encoded msg to peer", conn, encoded = shortLog(msg)
|
|
|
|
let f = conn.writeLp(msg)
|
|
|
|
if not f.completed():
|
|
|
|
sendMsgContinue(conn, f)
|
|
|
|
else:
|
|
|
|
f
|
|
|
|
else:
|
|
|
|
sendMsgSlow(p, msg)
|
|
|
|
|
|
|
|
proc sendEncoded*(p: PubSubPeer, msg: seq[byte], isHighPriority: bool): Future[void] =
|
|
|
|
## Asynchronously sends an encoded message to a specified `PubSubPeer`.
|
|
|
|
##
|
|
|
|
## Parameters:
|
|
|
|
## - `p`: The `PubSubPeer` instance to which the message is to be sent.
|
|
|
|
## - `msg`: The message to be sent, encoded as a sequence of bytes (`seq[byte]`).
|
|
|
|
## - `isHighPriority`: A boolean indicating whether the message should be treated as high priority.
|
|
|
|
## High priority messages are sent immediately, while low priority messages are queued and sent only after all high
|
|
|
|
## priority messages have been sent.
|
|
|
|
doAssert(not isNil(p), "pubsubpeer nil!")
|
|
|
|
|
2024-05-02 12:26:16 +02:00
|
|
|
p.clearSendPriorityQueue()
|
|
|
|
|
|
|
|
# When queues are empty, skipping the non-priority queue for low priority
|
|
|
|
# messages reduces latency
|
|
|
|
let emptyQueues =
|
|
|
|
(p.rpcmessagequeue.sendPriorityQueue.len() +
|
|
|
|
p.rpcmessagequeue.nonPriorityQueue.len()) == 0
|
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
if msg.len <= 0:
|
|
|
|
debug "empty message, skipping", p, msg = shortLog(msg)
|
|
|
|
Future[void].completed()
|
|
|
|
elif msg.len > p.maxMessageSize:
|
|
|
|
info "trying to send a msg too big for pubsub", maxSize=p.maxMessageSize, msgSize=msg.len
|
|
|
|
Future[void].completed()
|
2024-05-02 12:26:16 +02:00
|
|
|
elif isHighPriority or emptyQueues:
|
2024-03-05 16:05:21 +01:00
|
|
|
let f = p.sendMsg(msg)
|
|
|
|
if not f.finished:
|
|
|
|
p.rpcmessagequeue.sendPriorityQueue.addLast(f)
|
|
|
|
when defined(pubsubpeer_queue_metrics):
|
|
|
|
libp2p_gossipsub_priority_queue_size.inc(labelValues = [$p.peerId])
|
|
|
|
f
|
|
|
|
else:
|
2024-03-25 22:00:11 +01:00
|
|
|
if len(p.rpcmessagequeue.nonPriorityQueue) >= p.maxNumElementsInNonPriorityQueue:
|
|
|
|
if not p.disconnected:
|
|
|
|
p.disconnected = true
|
|
|
|
libp2p_pubsub_disconnects_over_non_priority_queue_limit.inc()
|
|
|
|
p.closeSendConn(PubSubPeerEventKind.DisconnectionRequested)
|
|
|
|
else:
|
|
|
|
Future[void].completed()
|
|
|
|
else:
|
|
|
|
let f = p.rpcmessagequeue.nonPriorityQueue.addLast(msg)
|
|
|
|
when defined(pubsubpeer_queue_metrics):
|
|
|
|
libp2p_gossipsub_non_priority_queue_size.inc(labelValues = [$p.peerId])
|
|
|
|
f
|
2021-04-18 10:08:33 +02:00
|
|
|
|
2023-10-02 11:39:28 +02:00
|
|
|
iterator splitRPCMsg(peer: PubSubPeer, rpcMsg: RPCMsg, maxSize: int, anonymize: bool): seq[byte] =
|
|
|
|
## This iterator takes an `RPCMsg` and sequentially repackages its Messages into new `RPCMsg` instances.
|
|
|
|
## Each new `RPCMsg` accumulates Messages until reaching the specified `maxSize`. If a single Message
|
|
|
|
## exceeds the `maxSize` when trying to fit into an empty `RPCMsg`, the latter is skipped as too large to send.
|
|
|
|
## Every constructed `RPCMsg` is then encoded, optionally anonymized, and yielded as a sequence of bytes.
|
|
|
|
|
|
|
|
var currentRPCMsg = rpcMsg
|
|
|
|
currentRPCMsg.messages = newSeq[Message]()
|
|
|
|
|
|
|
|
var currentSize = byteSize(currentRPCMsg)
|
|
|
|
|
|
|
|
for msg in rpcMsg.messages:
|
|
|
|
let msgSize = byteSize(msg)
|
|
|
|
|
|
|
|
# Check if adding the next message will exceed maxSize
|
|
|
|
if float(currentSize + msgSize) * 1.1 > float(maxSize): # Guessing 10% protobuf overhead
|
|
|
|
if currentRPCMsg.messages.len == 0:
|
|
|
|
trace "message too big to sent", peer, rpcMsg = shortLog(currentRPCMsg)
|
|
|
|
continue # Skip this message
|
|
|
|
|
|
|
|
trace "sending msg to peer", peer, rpcMsg = shortLog(currentRPCMsg)
|
|
|
|
yield encodeRpcMsg(currentRPCMsg, anonymize)
|
|
|
|
currentRPCMsg = RPCMsg()
|
|
|
|
currentSize = 0
|
2020-09-24 18:43:20 +02:00
|
|
|
|
2023-10-02 11:39:28 +02:00
|
|
|
currentRPCMsg.messages.add(msg)
|
|
|
|
currentSize += msgSize
|
|
|
|
|
|
|
|
# Check if there is a non-empty currentRPCMsg left to be added
|
|
|
|
if currentSize > 0 and currentRPCMsg.messages.len > 0:
|
|
|
|
trace "sending msg to peer", peer, rpcMsg = shortLog(currentRPCMsg)
|
|
|
|
yield encodeRpcMsg(currentRPCMsg, anonymize)
|
|
|
|
else:
|
|
|
|
trace "message too big to sent", peer, rpcMsg = shortLog(currentRPCMsg)
|
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
proc send*(p: PubSubPeer, msg: RPCMsg, anonymize: bool, isHighPriority: bool) {.raises: [].} =
|
|
|
|
## Asynchronously sends an `RPCMsg` to a specified `PubSubPeer` with an option for anonymization.
|
|
|
|
##
|
|
|
|
## Parameters:
|
|
|
|
## - `p`: The `PubSubPeer` instance to which the message is to be sent.
|
|
|
|
## - `msg`: The `RPCMsg` instance representing the message to be sent.
|
|
|
|
## - `anonymize`: A boolean flag indicating whether the message should be sent with anonymization.
|
|
|
|
## - `isHighPriority`: A boolean flag indicating whether the message should be treated as high priority.
|
|
|
|
## High priority messages are sent immediately, while low priority messages are queued and sent only after all high
|
|
|
|
## priority messages have been sent.
|
2020-09-25 18:39:34 +02:00
|
|
|
# When sending messages, we take care to re-encode them with the right
|
|
|
|
# anonymization flag to ensure that we're not penalized for sending invalid
|
|
|
|
# or malicious data on the wire - in particular, re-encoding protects against
|
|
|
|
# some forms of valid but redundantly encoded protobufs with unknown or
|
|
|
|
# duplicated fields
|
2020-09-24 18:43:20 +02:00
|
|
|
let encoded = if p.hasObservers():
|
2020-11-23 15:02:23 -06:00
|
|
|
var mm = msg
|
2020-09-24 18:43:20 +02:00
|
|
|
# trigger send hooks
|
|
|
|
p.sendObservers(mm)
|
2020-11-23 15:02:23 -06:00
|
|
|
sendMetrics(mm)
|
2020-09-25 18:39:34 +02:00
|
|
|
encodeRpcMsg(mm, anonymize)
|
2020-09-24 18:43:20 +02:00
|
|
|
else:
|
|
|
|
# If there are no send hooks, we redundantly re-encode the message to
|
|
|
|
# protobuf for every peer - this could easily be improved!
|
2020-11-23 15:02:23 -06:00
|
|
|
sendMetrics(msg)
|
2020-09-25 18:39:34 +02:00
|
|
|
encodeRpcMsg(msg, anonymize)
|
2020-09-24 18:43:20 +02:00
|
|
|
|
2023-10-02 11:39:28 +02:00
|
|
|
if encoded.len > p.maxMessageSize and msg.messages.len > 1:
|
|
|
|
for encodedSplitMsg in splitRPCMsg(p, msg, p.maxMessageSize, anonymize):
|
2024-03-05 16:05:21 +01:00
|
|
|
asyncSpawn p.sendEncoded(encodedSplitMsg, isHighPriority)
|
2023-10-02 11:39:28 +02:00
|
|
|
else:
|
|
|
|
# If the message size is within limits, send it as is
|
|
|
|
trace "sending msg to peer", peer = p, rpcMsg = shortLog(msg)
|
2024-03-05 16:05:21 +01:00
|
|
|
asyncSpawn p.sendEncoded(encoded, isHighPriority)
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2023-04-03 10:56:20 +02:00
|
|
|
proc canAskIWant*(p: PubSubPeer, msgId: MessageId): bool =
|
|
|
|
for sentIHave in p.sentIHaves.mitems():
|
|
|
|
if msgId in sentIHave:
|
|
|
|
sentIHave.excl(msgId)
|
|
|
|
return true
|
|
|
|
return false
|
|
|
|
|
2024-03-05 16:05:21 +01:00
|
|
|
proc sendNonPriorityTask(p: PubSubPeer) {.async.} =
|
|
|
|
while true:
|
|
|
|
# we send non-priority messages only if there are no pending priority messages
|
|
|
|
let msg = await p.rpcmessagequeue.nonPriorityQueue.popFirst()
|
|
|
|
while p.rpcmessagequeue.sendPriorityQueue.len > 0:
|
|
|
|
p.clearSendPriorityQueue()
|
|
|
|
# waiting for the last future minimizes the number of times we have to
|
|
|
|
# wait for something (each wait = performance cost) -
|
|
|
|
# clearSendPriorityQueue ensures we're not waiting for an already-finished
|
|
|
|
# future
|
|
|
|
if p.rpcmessagequeue.sendPriorityQueue.len > 0:
|
2024-03-20 11:54:32 +01:00
|
|
|
# `race` prevents `p.rpcmessagequeue.sendPriorityQueue[^1]` from being
|
|
|
|
# cancelled when this task is cancelled
|
|
|
|
discard await race(p.rpcmessagequeue.sendPriorityQueue[^1])
|
2024-03-05 16:05:21 +01:00
|
|
|
when defined(pubsubpeer_queue_metrics):
|
|
|
|
libp2p_gossipsub_non_priority_queue_size.dec(labelValues = [$p.peerId])
|
|
|
|
await p.sendMsg(msg)
|
|
|
|
|
|
|
|
proc startSendNonPriorityTask(p: PubSubPeer) =
|
|
|
|
debug "starting sendNonPriorityTask", p
|
|
|
|
if p.rpcmessagequeue.sendNonPriorityTask.isNil:
|
|
|
|
p.rpcmessagequeue.sendNonPriorityTask = p.sendNonPriorityTask()
|
|
|
|
|
|
|
|
proc stopSendNonPriorityTask*(p: PubSubPeer) =
|
|
|
|
if not p.rpcmessagequeue.sendNonPriorityTask.isNil:
|
|
|
|
debug "stopping sendNonPriorityTask", p
|
|
|
|
p.rpcmessagequeue.sendNonPriorityTask.cancelSoon()
|
|
|
|
p.rpcmessagequeue.sendNonPriorityTask = nil
|
|
|
|
p.rpcmessagequeue.sendPriorityQueue.clear()
|
|
|
|
p.rpcmessagequeue.nonPriorityQueue.clear()
|
|
|
|
|
|
|
|
when defined(pubsubpeer_queue_metrics):
|
|
|
|
libp2p_gossipsub_priority_queue_size.set(labelValues = [$p.peerId], value = 0)
|
|
|
|
libp2p_gossipsub_non_priority_queue_size.set(labelValues = [$p.peerId], value = 0)
|
|
|
|
|
|
|
|
proc new(T: typedesc[RpcMessageQueue]): T =
|
|
|
|
return T(
|
|
|
|
sendPriorityQueue: initDeque[Future[void]](),
|
2024-03-25 22:00:11 +01:00
|
|
|
nonPriorityQueue: newAsyncQueue[seq[byte]]()
|
2024-03-05 16:05:21 +01:00
|
|
|
)
|
|
|
|
|
2021-06-07 09:32:08 +02:00
|
|
|
proc new*(
|
|
|
|
T: typedesc[PubSubPeer],
|
2021-12-16 11:05:20 +01:00
|
|
|
peerId: PeerId,
|
2021-06-07 09:32:08 +02:00
|
|
|
getConn: GetConn,
|
|
|
|
onEvent: OnEvent,
|
2021-10-25 12:58:38 +02:00
|
|
|
codec: string,
|
2023-09-22 16:45:08 +02:00
|
|
|
maxMessageSize: int,
|
2024-03-25 22:00:11 +01:00
|
|
|
maxNumElementsInNonPriorityQueue: int = DefaultMaxNumElementsInNonPriorityQueue,
|
2023-09-22 16:45:08 +02:00
|
|
|
overheadRateLimitOpt: Opt[TokenBucket] = Opt.none(TokenBucket)): T =
|
2021-06-07 09:32:08 +02:00
|
|
|
|
2023-04-03 10:56:20 +02:00
|
|
|
result = T(
|
2020-09-22 09:05:53 +02:00
|
|
|
getConn: getConn,
|
|
|
|
onEvent: onEvent,
|
|
|
|
codec: codec,
|
|
|
|
peerId: peerId,
|
2023-01-10 13:33:14 +01:00
|
|
|
connectedFut: newFuture[void](),
|
2023-09-22 16:45:08 +02:00
|
|
|
maxMessageSize: maxMessageSize,
|
2024-03-05 16:05:21 +01:00
|
|
|
overheadRateLimitOpt: overheadRateLimitOpt,
|
|
|
|
rpcmessagequeue: RpcMessageQueue.new(),
|
2024-03-25 22:00:11 +01:00
|
|
|
maxNumElementsInNonPriorityQueue: maxNumElementsInNonPriorityQueue
|
2020-09-22 09:05:53 +02:00
|
|
|
)
|
2023-04-03 10:56:20 +02:00
|
|
|
result.sentIHaves.addFirst(default(HashSet[MessageId]))
|
2024-05-02 12:18:55 +02:00
|
|
|
result.heDontWants.addFirst(default(HashSet[SaltedId]))
|
2024-03-05 16:05:21 +01:00
|
|
|
result.startSendNonPriorityTask()
|