2019-12-05 20:16:18 -06:00
|
|
|
## Nim-LibP2P
|
|
|
|
## Copyright (c) 2019 Status Research & Development GmbH
|
|
|
|
## 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.
|
|
|
|
|
2021-05-21 10:27:01 -06:00
|
|
|
{.push raises: [Defect].}
|
|
|
|
|
2021-03-09 13:22:52 +01:00
|
|
|
import std/[tables, sets, options, sequtils, random]
|
2021-04-18 10:08:33 +02:00
|
|
|
import chronos, chronicles, metrics
|
2020-09-04 08:10:32 +02:00
|
|
|
import ./pubsub,
|
|
|
|
./floodsub,
|
|
|
|
./pubsubpeer,
|
|
|
|
./peertable,
|
|
|
|
./mcache,
|
|
|
|
./timedcache,
|
|
|
|
./rpc/[messages, message],
|
2019-12-05 20:16:18 -06:00
|
|
|
../protocol,
|
2020-06-19 11:29:43 -06:00
|
|
|
../../stream/connection,
|
2020-09-04 08:10:32 +02:00
|
|
|
../../peerinfo,
|
2020-07-01 15:25:09 +09:00
|
|
|
../../peerid,
|
2021-01-13 23:49:44 +09:00
|
|
|
../../utility,
|
2021-01-15 13:48:03 +09:00
|
|
|
../../switch
|
2021-02-06 09:13:04 +09:00
|
|
|
|
2020-09-21 18:16:29 +09:00
|
|
|
import stew/results
|
|
|
|
export results
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2021-05-21 10:27:01 -06:00
|
|
|
import ./gossipsub/[types, scoring, behavior]
|
|
|
|
|
|
|
|
export types, scoring, behavior, pubsub
|
2021-02-06 09:13:04 +09:00
|
|
|
|
2019-12-05 20:16:18 -06:00
|
|
|
logScope:
|
2020-12-01 11:34:27 -06:00
|
|
|
topics = "libp2p gossipsub"
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2021-01-08 14:21:24 +09:00
|
|
|
declareCounter(libp2p_gossipsub_failed_publish, "number of failed publish")
|
2021-01-15 18:55:54 +09:00
|
|
|
declareCounter(libp2p_gossipsub_invalid_topic_subscription, "number of invalid topic subscriptions that happened")
|
2020-11-04 23:18:00 +09:00
|
|
|
|
2020-09-21 18:16:29 +09:00
|
|
|
proc init*(_: type[GossipSubParams]): GossipSubParams =
|
|
|
|
GossipSubParams(
|
|
|
|
explicit: true,
|
|
|
|
pruneBackoff: 1.minutes,
|
|
|
|
floodPublish: true,
|
|
|
|
gossipFactor: 0.25,
|
2020-11-19 16:48:17 +09:00
|
|
|
d: GossipSubD,
|
|
|
|
dLow: GossipSubDlo,
|
|
|
|
dHigh: GossipSubDhi,
|
|
|
|
dScore: GossipSubDlo,
|
|
|
|
dOut: GossipSubDlo - 1, # DLow - 1
|
|
|
|
dLazy: GossipSubD, # Like D
|
|
|
|
heartbeatInterval: GossipSubHeartbeatInterval,
|
|
|
|
historyLength: GossipSubHistoryLength,
|
|
|
|
historyGossip: GossipSubHistoryGossip,
|
|
|
|
fanoutTTL: GossipSubFanoutTTL,
|
2020-11-26 14:45:10 +09:00
|
|
|
seenTTL: 2.minutes,
|
2021-02-12 12:27:26 +09:00
|
|
|
gossipThreshold: -100,
|
|
|
|
publishThreshold: -1000,
|
2020-09-21 18:16:29 +09:00
|
|
|
graylistThreshold: -10000,
|
2020-10-03 09:26:45 +09:00
|
|
|
opportunisticGraftThreshold: 0,
|
2020-09-21 18:16:29 +09:00
|
|
|
decayInterval: 1.seconds,
|
|
|
|
decayToZero: 0.01,
|
2021-02-12 12:27:26 +09:00
|
|
|
retainScore: 2.minutes,
|
2020-09-21 18:16:29 +09:00
|
|
|
appSpecificWeight: 0.0,
|
|
|
|
ipColocationFactorWeight: 0.0,
|
|
|
|
ipColocationFactorThreshold: 1.0,
|
|
|
|
behaviourPenaltyWeight: -1.0,
|
|
|
|
behaviourPenaltyDecay: 0.999,
|
2021-01-15 13:48:03 +09:00
|
|
|
disconnectBadPeers: false
|
2020-09-21 18:16:29 +09:00
|
|
|
)
|
|
|
|
|
|
|
|
proc validateParameters*(parameters: GossipSubParams): Result[void, cstring] =
|
2020-11-19 16:48:17 +09:00
|
|
|
if (parameters.dOut >= parameters.dLow) or
|
|
|
|
(parameters.dOut > (parameters.d div 2)):
|
2020-09-21 18:16:29 +09:00
|
|
|
err("gossipsub: dOut parameter error, Number of outbound connections to keep in the mesh. Must be less than D_lo and at most D/2")
|
|
|
|
elif parameters.gossipThreshold >= 0:
|
|
|
|
err("gossipsub: gossipThreshold parameter error, Must be < 0")
|
|
|
|
elif parameters.publishThreshold >= parameters.gossipThreshold:
|
|
|
|
err("gossipsub: publishThreshold parameter error, Must be < gossipThreshold")
|
|
|
|
elif parameters.graylistThreshold >= parameters.publishThreshold:
|
|
|
|
err("gossipsub: graylistThreshold parameter error, Must be < publishThreshold")
|
|
|
|
elif parameters.acceptPXThreshold < 0:
|
|
|
|
err("gossipsub: acceptPXThreshold parameter error, Must be >= 0")
|
|
|
|
elif parameters.opportunisticGraftThreshold < 0:
|
|
|
|
err("gossipsub: opportunisticGraftThreshold parameter error, Must be >= 0")
|
|
|
|
elif parameters.decayToZero > 0.5 or parameters.decayToZero <= 0.0:
|
|
|
|
err("gossipsub: decayToZero parameter error, Should be close to 0.0")
|
|
|
|
elif parameters.appSpecificWeight < 0:
|
|
|
|
err("gossipsub: appSpecificWeight parameter error, Must be positive")
|
|
|
|
elif parameters.ipColocationFactorWeight > 0:
|
|
|
|
err("gossipsub: ipColocationFactorWeight parameter error, Must be negative or 0")
|
|
|
|
elif parameters.ipColocationFactorThreshold < 1.0:
|
|
|
|
err("gossipsub: ipColocationFactorThreshold parameter error, Must be at least 1")
|
|
|
|
elif parameters.behaviourPenaltyWeight >= 0:
|
|
|
|
err("gossipsub: behaviourPenaltyWeight parameter error, Must be negative")
|
|
|
|
elif parameters.behaviourPenaltyDecay < 0 or parameters.behaviourPenaltyDecay >= 1:
|
|
|
|
err("gossipsub: behaviourPenaltyDecay parameter error, Must be between 0 and 1")
|
|
|
|
else:
|
|
|
|
ok()
|
|
|
|
|
|
|
|
proc validateParameters*(parameters: TopicParams): Result[void, cstring] =
|
|
|
|
if parameters.timeInMeshWeight <= 0.0 or parameters.timeInMeshWeight > 1.0:
|
|
|
|
err("gossipsub: timeInMeshWeight parameter error, Must be a small positive value")
|
|
|
|
elif parameters.timeInMeshCap <= 0.0:
|
|
|
|
err("gossipsub: timeInMeshCap parameter error, Should be a positive value")
|
|
|
|
elif parameters.firstMessageDeliveriesWeight <= 0.0:
|
|
|
|
err("gossipsub: firstMessageDeliveriesWeight parameter error, Should be a positive value")
|
|
|
|
elif parameters.meshMessageDeliveriesWeight >= 0.0:
|
|
|
|
err("gossipsub: meshMessageDeliveriesWeight parameter error, Should be a negative value")
|
|
|
|
elif parameters.meshMessageDeliveriesThreshold <= 0.0:
|
|
|
|
err("gossipsub: meshMessageDeliveriesThreshold parameter error, Should be a positive value")
|
|
|
|
elif parameters.meshMessageDeliveriesCap < parameters.meshMessageDeliveriesThreshold:
|
|
|
|
err("gossipsub: meshMessageDeliveriesCap parameter error, Should be >= meshMessageDeliveriesThreshold")
|
|
|
|
elif parameters.meshFailurePenaltyWeight >= 0.0:
|
|
|
|
err("gossipsub: meshFailurePenaltyWeight parameter error, Should be a negative value")
|
|
|
|
elif parameters.invalidMessageDeliveriesWeight >= 0.0:
|
|
|
|
err("gossipsub: invalidMessageDeliveriesWeight parameter error, Should be a negative value")
|
|
|
|
else:
|
|
|
|
ok()
|
|
|
|
|
2020-06-13 07:54:12 +08:00
|
|
|
method init*(g: GossipSub) =
|
2019-12-16 23:24:03 -06:00
|
|
|
proc handler(conn: Connection, proto: string) {.async.} =
|
2019-12-05 20:16:18 -06:00
|
|
|
## main protocol handler that gets triggered on every
|
|
|
|
## connection for a protocol string
|
|
|
|
## e.g. ``/floodsub/1.0.0``, etc...
|
|
|
|
##
|
2020-09-04 19:30:45 +03:00
|
|
|
try:
|
|
|
|
await g.handleConn(conn, proto)
|
|
|
|
except CancelledError:
|
|
|
|
# This is top-level procedure which will work as separate task, so it
|
|
|
|
# do not need to propogate CancelledError.
|
2020-09-06 10:31:47 +02:00
|
|
|
trace "Unexpected cancellation in gossipsub handler", conn
|
2020-09-04 19:30:45 +03:00
|
|
|
except CatchableError as exc:
|
2020-09-06 10:31:47 +02:00
|
|
|
trace "GossipSub handler leaks an error", exc = exc.msg, conn
|
2019-12-05 20:16:18 -06:00
|
|
|
|
|
|
|
g.handler = handler
|
2020-09-21 18:16:29 +09:00
|
|
|
g.codecs &= GossipSubCodec
|
|
|
|
g.codecs &= GossipSubCodec_10
|
|
|
|
|
|
|
|
method onNewPeer(g: GossipSub, peer: PubSubPeer) =
|
2021-03-09 13:22:52 +01:00
|
|
|
g.withPeerStats(peer.peerId) do (stats: var PeerStats):
|
|
|
|
# Make sure stats and peer information match, even when reloading peer stats
|
|
|
|
# from a previous connection
|
2021-01-13 23:49:44 +09:00
|
|
|
peer.score = stats.score
|
|
|
|
peer.appScore = stats.appScore
|
|
|
|
peer.behaviourPenalty = stats.behaviourPenalty
|
2020-09-21 18:16:29 +09:00
|
|
|
|
2021-03-09 13:22:52 +01:00
|
|
|
peer.iWantBudget = IWantPeerBudget
|
|
|
|
peer.iHaveBudget = IHavePeerBudget
|
|
|
|
|
2020-09-22 09:05:53 +02:00
|
|
|
method onPubSubPeerEvent*(p: GossipSub, peer: PubsubPeer, event: PubSubPeerEvent) {.gcsafe.} =
|
|
|
|
case event.kind
|
|
|
|
of PubSubPeerEventKind.Connected:
|
|
|
|
discard
|
|
|
|
of PubSubPeerEventKind.Disconnected:
|
|
|
|
# If a send connection is lost, it's better to remove peer from the mesh -
|
|
|
|
# if it gets reestablished, the peer will be readded to the mesh, and if it
|
|
|
|
# doesn't, well.. then we hope the peer is going away!
|
2020-10-02 13:09:31 +09:00
|
|
|
for topic, peers in p.mesh.mpairs():
|
|
|
|
p.pruned(peer, topic)
|
2020-09-22 09:05:53 +02:00
|
|
|
peers.excl(peer)
|
|
|
|
for _, peers in p.fanout.mpairs():
|
|
|
|
peers.excl(peer)
|
|
|
|
|
|
|
|
procCall FloodSub(p).onPubSubPeerEvent(peer, event)
|
|
|
|
|
2020-08-11 18:05:49 -06:00
|
|
|
method unsubscribePeer*(g: GossipSub, peer: PeerID) =
|
2019-12-05 20:16:18 -06:00
|
|
|
## handle peer disconnects
|
2020-08-11 18:05:49 -06:00
|
|
|
##
|
|
|
|
|
2020-09-06 10:31:47 +02:00
|
|
|
trace "unsubscribing gossipsub peer", peer
|
2020-08-11 18:05:49 -06:00
|
|
|
let pubSubPeer = g.peers.getOrDefault(peer)
|
|
|
|
if pubSubPeer.isNil:
|
2020-09-06 10:31:47 +02:00
|
|
|
trace "no peer to unsubscribe", peer
|
2020-08-11 18:05:49 -06:00
|
|
|
return
|
2020-06-19 15:19:07 -06:00
|
|
|
|
2020-09-21 18:16:29 +09:00
|
|
|
# remove from peer IPs collection too
|
2021-02-27 23:49:56 +09:00
|
|
|
if pubSubPeer.address.isSome():
|
|
|
|
g.peersInIP.withValue(pubSubPeer.address.get(), s):
|
2021-02-27 21:31:59 +09:00
|
|
|
s[].excl(pubSubPeer.peerId)
|
2021-01-15 13:48:03 +09:00
|
|
|
if s[].len == 0:
|
2021-02-27 23:49:56 +09:00
|
|
|
g.peersInIP.del(pubSubPeer.address.get())
|
2020-09-21 18:16:29 +09:00
|
|
|
|
2020-08-11 18:05:49 -06:00
|
|
|
for t in toSeq(g.gossipsub.keys):
|
|
|
|
g.gossipsub.removePeer(t, pubSubPeer)
|
2020-09-21 18:16:29 +09:00
|
|
|
# also try to remove from explicit table here
|
|
|
|
g.explicit.removePeer(t, pubSubPeer)
|
2020-06-11 20:20:58 -06:00
|
|
|
|
2020-08-11 18:05:49 -06:00
|
|
|
for t in toSeq(g.mesh.keys):
|
2020-12-15 10:25:22 +09:00
|
|
|
trace "pruning unsubscribing peer", pubSubPeer, score = pubSubPeer.score
|
2020-09-21 18:16:29 +09:00
|
|
|
g.pruned(pubSubPeer, t)
|
2020-08-11 18:05:49 -06:00
|
|
|
g.mesh.removePeer(t, pubSubPeer)
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2020-08-11 18:05:49 -06:00
|
|
|
for t in toSeq(g.fanout.keys):
|
|
|
|
g.fanout.removePeer(t, pubSubPeer)
|
2020-08-06 11:12:52 +09:00
|
|
|
|
2021-01-08 14:21:24 +09:00
|
|
|
g.peerStats.withValue(peer, stats):
|
2020-12-15 10:25:22 +09:00
|
|
|
for topic, info in stats[].topicInfos.mpairs:
|
2020-09-21 18:16:29 +09:00
|
|
|
info.firstMessageDeliveries = 0
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2020-08-11 18:05:49 -06:00
|
|
|
procCall FloodSub(g).unsubscribePeer(peer)
|
2020-05-21 11:33:48 -06:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
proc handleSubscribe*(g: GossipSub,
|
|
|
|
peer: PubSubPeer,
|
|
|
|
topic: string,
|
|
|
|
subscribe: bool) =
|
2020-07-16 21:26:57 +02:00
|
|
|
logScope:
|
2020-09-06 10:31:47 +02:00
|
|
|
peer
|
2020-07-16 21:26:57 +02:00
|
|
|
topic
|
2020-09-22 09:05:53 +02:00
|
|
|
|
2019-12-10 14:50:35 -06:00
|
|
|
if subscribe:
|
2021-05-07 00:43:45 +02:00
|
|
|
# this is a workaround for a race condition
|
|
|
|
# that can happen if we disconnect the peer very early
|
|
|
|
# in the future we might use this as a test case
|
|
|
|
# and eventually remove this workaround
|
|
|
|
if peer.peerId notin g.peers:
|
|
|
|
trace "ignoring unknown peer"
|
|
|
|
return
|
|
|
|
|
|
|
|
if not(isNil(g.subscriptionValidator)) and not(g.subscriptionValidator(topic)):
|
|
|
|
# this is a violation, so warn should be in order
|
|
|
|
trace "ignoring invalid topic subscription", topic, peer
|
|
|
|
libp2p_gossipsub_invalid_topic_subscription.inc()
|
|
|
|
return
|
|
|
|
|
2020-07-16 21:26:57 +02:00
|
|
|
trace "peer subscribed to topic"
|
2021-01-13 23:49:44 +09:00
|
|
|
|
2020-05-27 12:33:49 -06:00
|
|
|
# subscribe remote peer to the topic
|
2020-07-13 22:32:38 +09:00
|
|
|
discard g.gossipsub.addPeer(topic, peer)
|
2020-09-21 18:16:29 +09:00
|
|
|
if peer.peerId in g.parameters.directPeers:
|
|
|
|
discard g.explicit.addPeer(topic, peer)
|
2019-12-10 14:50:35 -06:00
|
|
|
else:
|
2020-07-16 21:26:57 +02:00
|
|
|
trace "peer unsubscribed from topic"
|
2021-01-13 23:49:44 +09:00
|
|
|
|
2020-05-27 12:33:49 -06:00
|
|
|
# unsubscribe remote peer from the topic
|
2020-07-13 22:32:38 +09:00
|
|
|
g.gossipsub.removePeer(topic, peer)
|
|
|
|
g.mesh.removePeer(topic, peer)
|
|
|
|
g.fanout.removePeer(topic, peer)
|
2020-09-21 18:16:29 +09:00
|
|
|
if peer.peerId in g.parameters.directPeers:
|
|
|
|
g.explicit.removePeer(topic, peer)
|
2020-07-09 19:16:46 +02:00
|
|
|
|
|
|
|
trace "gossip peers", peers = g.gossipsub.peers(topic), topic
|
2020-07-07 18:33:05 -06:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
proc handleControl(g: GossipSub, peer: PubSubPeer, control: ControlMessage) =
|
|
|
|
g.handlePrune(peer, control.prune)
|
|
|
|
|
|
|
|
var respControl: ControlMessage
|
|
|
|
let iwant = g.handleIHave(peer, control.ihave)
|
|
|
|
if iwant.messageIDs.len > 0:
|
|
|
|
respControl.iwant.add(iwant)
|
|
|
|
respControl.prune.add(g.handleGraft(peer, control.graft))
|
|
|
|
let messages = g.handleIWant(peer, control.iwant)
|
|
|
|
|
|
|
|
if
|
|
|
|
respControl.prune.len > 0 or
|
|
|
|
respControl.iwant.len > 0 or
|
|
|
|
messages.len > 0:
|
|
|
|
# iwant and prunes from here, also messages
|
|
|
|
|
|
|
|
for smsg in messages:
|
|
|
|
for topic in smsg.topicIDs:
|
|
|
|
if g.knownTopics.contains(topic):
|
|
|
|
libp2p_pubsub_broadcast_messages.inc(labelValues = [topic])
|
2021-04-18 10:08:33 +02:00
|
|
|
else:
|
2021-05-07 00:43:45 +02:00
|
|
|
libp2p_pubsub_broadcast_messages.inc(labelValues = ["generic"])
|
2021-04-28 10:03:03 +09:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
libp2p_pubsub_broadcast_iwant.inc(respControl.iwant.len.int64)
|
|
|
|
|
|
|
|
for prune in respControl.prune:
|
|
|
|
if g.knownTopics.contains(prune.topicID):
|
|
|
|
libp2p_pubsub_broadcast_prune.inc(labelValues = [prune.topicID])
|
|
|
|
else:
|
|
|
|
libp2p_pubsub_broadcast_prune.inc(labelValues = ["generic"])
|
|
|
|
|
|
|
|
trace "sending control message", msg = shortLog(respControl), peer
|
|
|
|
g.send(
|
|
|
|
peer,
|
|
|
|
RPCMsg(control: some(respControl), messages: messages))
|
2021-04-18 10:08:33 +02:00
|
|
|
|
2020-06-13 07:54:12 +08:00
|
|
|
method rpcHandler*(g: GossipSub,
|
2019-12-05 20:16:18 -06:00
|
|
|
peer: PubSubPeer,
|
2020-09-01 09:33:03 +02:00
|
|
|
rpcMsg: RPCMsg) {.async.} =
|
2021-05-07 00:43:45 +02:00
|
|
|
for i in 0..<min(g.topicsHigh, rpcMsg.subscriptions.len):
|
|
|
|
template sub: untyped = rpcMsg.subscriptions[i]
|
|
|
|
g.handleSubscribe(peer, sub.topic, sub.subscribe)
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2021-02-12 12:27:26 +09:00
|
|
|
# the above call applied limtis to subs number
|
|
|
|
# in gossipsub we want to apply scoring as well
|
|
|
|
if rpcMsg.subscriptions.len > g.topicsHigh:
|
|
|
|
debug "received an rpc message with an oversized amount of subscriptions", peer,
|
|
|
|
size = rpcMsg.subscriptions.len,
|
|
|
|
limit = g.topicsHigh
|
|
|
|
peer.behaviourPenalty += 0.1
|
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
for i in 0..<rpcMsg.messages.len(): # for every message
|
|
|
|
template msg: untyped = rpcMsg.messages[i]
|
2020-09-04 08:10:32 +02:00
|
|
|
let msgId = g.msgIdProvider(msg)
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2021-01-13 23:49:44 +09:00
|
|
|
# avoid the remote peer from controlling the seen table hashing
|
|
|
|
# by adding random bytes to the ID we ensure we randomize the IDs
|
|
|
|
# we do only for seen as this is the great filter from the external world
|
2021-04-18 10:08:33 +02:00
|
|
|
if g.addSeen(msgId):
|
2020-12-18 16:14:07 +01:00
|
|
|
trace "Dropping already-seen message", msgId = shortLog(msgId), peer
|
2020-09-21 18:16:29 +09:00
|
|
|
# make sure to update score tho before continuing
|
2021-04-18 10:08:33 +02:00
|
|
|
# TODO: take into account meshMessageDeliveriesWindow
|
|
|
|
# score only if messages are not too old.
|
|
|
|
g.rewardDelivered(peer, msg.topicIDs, false)
|
2020-12-15 22:46:03 +01:00
|
|
|
|
2020-12-15 10:25:22 +09:00
|
|
|
# onto the next message
|
2020-09-04 08:10:32 +02:00
|
|
|
continue
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2021-01-13 23:49:44 +09:00
|
|
|
# avoid processing messages we are not interested in
|
|
|
|
if msg.topicIDs.allIt(it notin g.topics):
|
|
|
|
debug "Dropping message of topic without subscription", msgId = shortLog(msgId), peer
|
|
|
|
continue
|
|
|
|
|
2020-09-24 00:56:33 +09:00
|
|
|
if (msg.signature.len > 0 or g.verifySignature) and not msg.verify():
|
|
|
|
# always validate if signature is present or required
|
2020-12-18 16:14:07 +01:00
|
|
|
debug "Dropping message due to failed signature verification",
|
|
|
|
msgId = shortLog(msgId), peer
|
2021-01-25 21:13:42 +09:00
|
|
|
g.punishInvalidMessage(peer, msg.topicIDs)
|
2020-09-04 08:10:32 +02:00
|
|
|
continue
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2020-09-24 00:56:33 +09:00
|
|
|
if msg.seqno.len > 0 and msg.seqno.len != 8:
|
|
|
|
# if we have seqno should be 8 bytes long
|
2020-12-18 16:14:07 +01:00
|
|
|
debug "Dropping message due to invalid seqno length",
|
|
|
|
msgId = shortLog(msgId), peer
|
2021-01-25 21:13:42 +09:00
|
|
|
g.punishInvalidMessage(peer, msg.topicIDs)
|
2020-09-24 00:56:33 +09:00
|
|
|
continue
|
|
|
|
|
|
|
|
# g.anonymize needs no evaluation when receiving messages
|
|
|
|
# as we have a "lax" policy and allow signed messages
|
|
|
|
|
2020-10-21 12:25:42 +09:00
|
|
|
let validation = await g.validate(msg)
|
|
|
|
case validation
|
|
|
|
of ValidationResult.Reject:
|
2020-12-18 16:14:07 +01:00
|
|
|
debug "Dropping message after validation, reason: reject",
|
|
|
|
msgId = shortLog(msgId), peer
|
2021-01-25 21:13:42 +09:00
|
|
|
g.punishInvalidMessage(peer, msg.topicIDs)
|
2020-09-04 08:10:32 +02:00
|
|
|
continue
|
2020-10-21 12:25:42 +09:00
|
|
|
of ValidationResult.Ignore:
|
2020-12-18 16:14:07 +01:00
|
|
|
debug "Dropping message after validation, reason: ignore",
|
|
|
|
msgId = shortLog(msgId), peer
|
2020-10-21 12:25:42 +09:00
|
|
|
continue
|
|
|
|
of ValidationResult.Accept:
|
|
|
|
discard
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2020-12-15 10:25:22 +09:00
|
|
|
# store in cache only after validation
|
|
|
|
g.mcache.put(msgId, msg)
|
|
|
|
|
2021-04-18 10:08:33 +02:00
|
|
|
g.rewardDelivered(peer, msg.topicIDs, true)
|
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
var toSendPeers = initHashSet[PubSubPeer]()
|
2020-09-21 18:16:29 +09:00
|
|
|
for t in msg.topicIDs: # for every topic in the message
|
2021-01-13 23:49:44 +09:00
|
|
|
if t notin g.topics:
|
|
|
|
continue
|
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
g.floodsub.withValue(t, peers): toSendPeers.incl(peers[])
|
|
|
|
g.mesh.withValue(t, peers): toSendPeers.incl(peers[])
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
await handleData(g, t, msg.data)
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
# In theory, if topics are the same in all messages, we could batch - we'd
|
|
|
|
# also have to be careful to only include validated messages
|
2021-04-18 10:08:33 +02:00
|
|
|
g.broadcast(toSendPeers, RPCMsg(messages: @[msg]))
|
|
|
|
trace "forwared message to peers", peers = toSendPeers.len, msgId, peer
|
2021-01-08 14:21:24 +09:00
|
|
|
for topic in msg.topicIDs:
|
|
|
|
if g.knownTopics.contains(topic):
|
2021-04-18 10:08:33 +02:00
|
|
|
libp2p_pubsub_messages_rebroadcasted.inc(toSendPeers.len.int64, labelValues = [topic])
|
2021-01-08 14:21:24 +09:00
|
|
|
else:
|
2021-04-18 10:08:33 +02:00
|
|
|
libp2p_pubsub_messages_rebroadcasted.inc(toSendPeers.len.int64, labelValues = ["generic"])
|
2020-09-01 09:33:03 +02:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
if rpcMsg.control.isSome():
|
|
|
|
g.handleControl(peer, rpcMsg.control.unsafeGet())
|
2020-09-22 09:05:53 +02:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
g.updateMetrics(rpcMsg)
|
2020-09-22 09:05:53 +02:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
method onTopicSubscription*(g: GossipSub, topic: string, subscribed: bool) =
|
|
|
|
if subscribed:
|
|
|
|
procCall PubSub(g).onTopicSubscription(topic, subscribed)
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
# if we have a fanout on this topic break it
|
|
|
|
if topic in g.fanout:
|
|
|
|
g.fanout.del(topic)
|
2020-07-21 01:16:13 +09:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
# rebalance but don't update metrics here, we do that only in the heartbeat
|
|
|
|
g.rebalanceMesh(topic, metrics = nil)
|
|
|
|
else:
|
2020-12-19 23:43:32 +09:00
|
|
|
let mpeers = g.mesh.getOrDefault(topic)
|
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
# Remove peers from the mesh since we're no longer both interested
|
|
|
|
# in the topic
|
|
|
|
let msg = RPCMsg(control: some(ControlMessage(
|
|
|
|
prune: @[ControlPrune(
|
|
|
|
topicID: topic,
|
|
|
|
peers: g.peerExchangeList(topic),
|
|
|
|
backoff: g.parameters.pruneBackoff.seconds.uint64)])))
|
|
|
|
g.broadcast(mpeers, msg)
|
2020-12-19 23:43:32 +09:00
|
|
|
|
|
|
|
for peer in mpeers:
|
2020-09-21 18:16:29 +09:00
|
|
|
g.pruned(peer, topic)
|
2020-12-19 23:43:32 +09:00
|
|
|
|
2020-12-20 00:45:34 +09:00
|
|
|
g.mesh.del(topic)
|
|
|
|
|
2020-12-19 23:43:32 +09:00
|
|
|
|
2021-05-07 00:43:45 +02:00
|
|
|
# Send unsubscribe (in reverse order to sub/graft)
|
|
|
|
procCall PubSub(g).onTopicSubscription(topic, subscribed)
|
2019-12-05 20:16:18 -06:00
|
|
|
|
|
|
|
method publish*(g: GossipSub,
|
|
|
|
topic: string,
|
2020-09-01 09:33:03 +02:00
|
|
|
data: seq[byte]): Future[int] {.async.} =
|
2020-07-07 18:33:05 -06:00
|
|
|
# base returns always 0
|
2020-09-01 09:33:03 +02:00
|
|
|
discard await procCall PubSub(g).publish(topic, data)
|
2020-09-04 08:10:32 +02:00
|
|
|
|
2021-01-13 23:49:44 +09:00
|
|
|
logScope:
|
|
|
|
topic
|
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
trace "Publishing message on topic", data = data.shortLog
|
2021-05-24 11:55:33 -06:00
|
|
|
|
|
|
|
if topic.len <= 0: # data could be 0/empty
|
2020-09-04 08:10:32 +02:00
|
|
|
debug "Empty topic, skipping publish"
|
2020-07-07 18:33:05 -06:00
|
|
|
return 0
|
2020-09-22 09:05:53 +02:00
|
|
|
|
2021-05-24 11:55:33 -06:00
|
|
|
var peers: HashSet[PubSubPeer]
|
|
|
|
|
|
|
|
if g.parameters.floodPublish:
|
|
|
|
# With flood publishing enabled, the mesh is used when propagating messages from other peers,
|
|
|
|
# but a peer's own messages will always be published to all known peers in the topic.
|
|
|
|
for peer in g.gossipsub.getOrDefault(topic):
|
2020-09-21 18:16:29 +09:00
|
|
|
if peer.score >= g.parameters.publishThreshold:
|
2021-05-24 11:55:33 -06:00
|
|
|
trace "publish: including flood/high score peer", peer
|
|
|
|
peers.incl(peer)
|
|
|
|
|
|
|
|
# add always direct peers
|
|
|
|
peers.incl(g.explicit.getOrDefault(topic))
|
|
|
|
|
2020-07-07 18:33:05 -06:00
|
|
|
if topic in g.topics: # if we're subscribed use the mesh
|
2020-09-21 18:16:29 +09:00
|
|
|
peers.incl(g.mesh.getOrDefault(topic))
|
2020-07-07 18:33:05 -06:00
|
|
|
else: # not subscribed, send to fanout peers
|
|
|
|
# try optimistically
|
2020-09-21 18:16:29 +09:00
|
|
|
peers.incl(g.fanout.getOrDefault(topic))
|
2020-07-07 18:33:05 -06:00
|
|
|
if peers.len == 0:
|
|
|
|
# ok we had nothing.. let's try replenish inline
|
|
|
|
g.replenishFanout(topic)
|
2020-09-21 18:16:29 +09:00
|
|
|
peers.incl(g.fanout.getOrDefault(topic))
|
2020-07-07 18:33:05 -06:00
|
|
|
|
2020-07-09 14:21:47 -06:00
|
|
|
# even if we couldn't publish,
|
|
|
|
# we still attempted to publish
|
|
|
|
# on the topic, so it makes sense
|
|
|
|
# to update the last topic publish
|
|
|
|
# time
|
2020-11-19 16:48:17 +09:00
|
|
|
g.lastFanoutPubSub[topic] = Moment.fromNow(g.parameters.fanoutTTL)
|
2020-07-09 14:21:47 -06:00
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
if peers.len == 0:
|
2021-02-09 18:42:59 +09:00
|
|
|
let topicPeers = g.gossipsub.getOrDefault(topic).toSeq()
|
|
|
|
notice "No peers for topic, skipping publish", peersOnTopic = topicPeers.len,
|
|
|
|
connectedPeers = topicPeers.filterIt(it.connected).len,
|
|
|
|
topic
|
2021-01-08 14:21:24 +09:00
|
|
|
# skipping topic as our metrics finds that heavy
|
|
|
|
libp2p_gossipsub_failed_publish.inc()
|
2020-09-04 08:10:32 +02:00
|
|
|
return 0
|
|
|
|
|
2020-07-15 12:51:33 +09:00
|
|
|
inc g.msgSeqno
|
2020-07-07 18:33:05 -06:00
|
|
|
let
|
2020-09-24 00:56:33 +09:00
|
|
|
msg =
|
|
|
|
if g.anonymize:
|
|
|
|
Message.init(none(PeerInfo), data, topic, none(uint64), false)
|
|
|
|
else:
|
|
|
|
Message.init(some(g.peerInfo), data, topic, some(g.msgSeqno), g.sign)
|
2020-07-07 18:33:05 -06:00
|
|
|
msgId = g.msgIdProvider(msg)
|
|
|
|
|
2020-12-18 16:14:07 +01:00
|
|
|
logScope: msgId = shortLog(msgId)
|
2020-07-07 18:33:05 -06:00
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
trace "Created new message", msg = shortLog(msg), peers = peers.len
|
2020-06-16 22:14:02 -06:00
|
|
|
|
2021-04-18 10:08:33 +02:00
|
|
|
if g.addSeen(msgId):
|
2020-09-04 08:10:32 +02:00
|
|
|
# custom msgid providers might cause this
|
|
|
|
trace "Dropping already-seen message"
|
2020-08-15 21:50:31 +02:00
|
|
|
return 0
|
2020-07-07 18:33:05 -06:00
|
|
|
|
2020-09-04 08:10:32 +02:00
|
|
|
g.mcache.put(msgId, msg)
|
|
|
|
|
2021-04-18 10:08:33 +02:00
|
|
|
g.broadcast(peers, RPCMsg(messages: @[msg]))
|
|
|
|
|
2021-01-08 14:21:24 +09:00
|
|
|
if g.knownTopics.contains(topic):
|
2021-04-18 10:08:33 +02:00
|
|
|
libp2p_pubsub_messages_published.inc(peers.len.int64, labelValues = [topic])
|
2020-11-26 16:20:34 +09:00
|
|
|
else:
|
2021-04-18 10:08:33 +02:00
|
|
|
libp2p_pubsub_messages_published.inc(peers.len.int64, labelValues = ["generic"])
|
2020-09-04 08:10:32 +02:00
|
|
|
|
2021-04-22 18:51:22 +09:00
|
|
|
trace "Published message to peers", peers=peers.len
|
2020-09-04 08:10:32 +02:00
|
|
|
|
|
|
|
return peers.len
|
|
|
|
|
2020-09-21 18:16:29 +09:00
|
|
|
proc maintainDirectPeers(g: GossipSub) {.async.} =
|
|
|
|
while g.heartbeatRunning:
|
2021-01-15 13:48:03 +09:00
|
|
|
for id, addrs in g.parameters.directPeers:
|
2020-09-21 18:16:29 +09:00
|
|
|
let peer = g.peers.getOrDefault(id)
|
2021-01-15 13:48:03 +09:00
|
|
|
if isNil(peer):
|
|
|
|
trace "Attempting to dial a direct peer", peer = id
|
|
|
|
try:
|
|
|
|
# dial, internally connection will be stored
|
|
|
|
let _ = await g.switch.dial(id, addrs, g.codecs)
|
|
|
|
# populate the peer after it's connected
|
|
|
|
discard g.getOrCreatePeer(id, g.codecs)
|
|
|
|
except CancelledError:
|
|
|
|
trace "Direct peer dial canceled"
|
|
|
|
raise
|
|
|
|
except CatchableError as exc:
|
|
|
|
debug "Direct peer error dialing", msg = exc.msg
|
2020-09-21 18:16:29 +09:00
|
|
|
|
|
|
|
await sleepAsync(1.minutes)
|
|
|
|
|
2019-12-05 20:16:18 -06:00
|
|
|
method start*(g: GossipSub) {.async.} =
|
2020-07-07 18:33:05 -06:00
|
|
|
trace "gossipsub start"
|
2020-06-28 17:56:38 +02:00
|
|
|
|
2020-09-04 18:31:43 +02:00
|
|
|
if not g.heartbeatFut.isNil:
|
|
|
|
warn "Starting gossipsub twice"
|
|
|
|
return
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2020-06-20 19:56:55 +09:00
|
|
|
g.heartbeatRunning = true
|
|
|
|
g.heartbeatFut = g.heartbeat()
|
2020-09-21 18:16:29 +09:00
|
|
|
g.directPeersLoop = g.maintainDirectPeers()
|
2020-06-20 19:56:55 +09:00
|
|
|
|
2019-12-05 20:16:18 -06:00
|
|
|
method stop*(g: GossipSub) {.async.} =
|
2020-07-07 18:33:05 -06:00
|
|
|
trace "gossipsub stop"
|
2020-09-04 18:31:43 +02:00
|
|
|
if g.heartbeatFut.isNil:
|
|
|
|
warn "Stopping gossipsub without starting it"
|
|
|
|
return
|
2019-12-05 20:16:18 -06:00
|
|
|
|
|
|
|
# stop heartbeat interval
|
2020-06-20 19:56:55 +09:00
|
|
|
g.heartbeatRunning = false
|
2020-09-21 18:16:29 +09:00
|
|
|
g.directPeersLoop.cancel()
|
2020-06-20 19:56:55 +09:00
|
|
|
if not g.heartbeatFut.finished:
|
2020-07-07 18:33:05 -06:00
|
|
|
trace "awaiting last heartbeat"
|
2020-06-20 19:56:55 +09:00
|
|
|
await g.heartbeatFut
|
2020-07-07 18:33:05 -06:00
|
|
|
trace "heartbeat stopped"
|
2020-09-04 18:31:43 +02:00
|
|
|
g.heartbeatFut = nil
|
2019-12-05 20:16:18 -06:00
|
|
|
|
2021-05-21 10:27:01 -06:00
|
|
|
method initPubSub*(g: GossipSub)
|
|
|
|
{.raises: [Defect, InitializationError].} =
|
2019-12-05 20:16:18 -06:00
|
|
|
procCall FloodSub(g).initPubSub()
|
|
|
|
|
2020-09-21 18:16:29 +09:00
|
|
|
if not g.parameters.explicit:
|
|
|
|
g.parameters = GossipSubParams.init()
|
2020-09-22 09:05:53 +02:00
|
|
|
|
2021-05-21 10:27:01 -06:00
|
|
|
let validationRes = g.parameters.validateParameters()
|
|
|
|
if validationRes.isErr:
|
|
|
|
raise newException(InitializationError, $validationRes.error)
|
2020-09-21 18:16:29 +09:00
|
|
|
|
2020-05-21 14:24:20 -06:00
|
|
|
randomize()
|
2020-11-26 14:45:10 +09:00
|
|
|
|
|
|
|
# init the floodsub stuff here, we customize timedcache in gossip!
|
|
|
|
g.seen = TimedCache[MessageID].init(g.parameters.seenTTL)
|
|
|
|
|
|
|
|
# init gossip stuff
|
2020-11-19 16:48:17 +09:00
|
|
|
g.mcache = MCache.init(g.parameters.historyGossip, g.parameters.historyLength)
|