2019-12-06 02:16:18 +00:00
|
|
|
## Nim-Libp2p
|
|
|
|
## Copyright (c) 2018 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.
|
|
|
|
|
2020-05-08 20:58:23 +00:00
|
|
|
{.used.}
|
|
|
|
|
2020-01-10 03:59:27 +00:00
|
|
|
import unittest, sequtils, options, tables, sets
|
2020-06-03 02:21:11 +00:00
|
|
|
import chronos, stew/byteutils
|
2020-04-21 01:24:42 +00:00
|
|
|
import chronicles
|
|
|
|
import utils, ../../libp2p/[errors,
|
2020-07-01 06:25:09 +00:00
|
|
|
peerid,
|
2019-12-06 02:16:18 +00:00
|
|
|
peerinfo,
|
2020-06-19 17:29:43 +00:00
|
|
|
stream/connection,
|
2020-08-03 05:20:11 +00:00
|
|
|
stream/bufferstream,
|
2019-12-06 02:16:18 +00:00
|
|
|
crypto/crypto,
|
|
|
|
protocols/pubsub/pubsub,
|
2020-12-19 14:43:32 +00:00
|
|
|
protocols/pubsub/gossipsub,
|
2020-08-03 05:20:11 +00:00
|
|
|
protocols/pubsub/pubsubpeer,
|
2020-07-15 19:18:55 +00:00
|
|
|
protocols/pubsub/peertable,
|
2019-12-17 05:24:03 +00:00
|
|
|
protocols/pubsub/rpc/messages]
|
2020-05-08 20:10:06 +00:00
|
|
|
import ../helpers
|
2020-04-21 01:24:42 +00:00
|
|
|
|
2021-04-18 08:08:33 +00:00
|
|
|
proc `$`(peer: PubSubPeer): string = shortLog(peer)
|
|
|
|
|
2020-04-21 01:24:42 +00:00
|
|
|
proc waitSub(sender, receiver: auto; key: string) {.async, gcsafe.} =
|
|
|
|
if sender == receiver:
|
|
|
|
return
|
|
|
|
# turn things deterministic
|
|
|
|
# this is for testing purposes only
|
|
|
|
# peers can be inside `mesh` and `fanout`, not just `gossipsub`
|
|
|
|
var ceil = 15
|
2020-08-12 00:05:49 +00:00
|
|
|
let fsub = GossipSub(sender)
|
2020-09-21 09:16:29 +00:00
|
|
|
let ev = newAsyncEvent()
|
|
|
|
fsub.heartbeatEvents.add(ev)
|
|
|
|
|
|
|
|
# await first heartbeat
|
|
|
|
await ev.wait()
|
|
|
|
ev.clear()
|
|
|
|
|
2020-04-21 01:24:42 +00:00
|
|
|
while (not fsub.gossipsub.hasKey(key) or
|
2020-08-12 00:05:49 +00:00
|
|
|
not fsub.gossipsub.hasPeerID(key, receiver.peerInfo.peerId)) and
|
2020-04-21 01:24:42 +00:00
|
|
|
(not fsub.mesh.hasKey(key) or
|
2020-08-12 00:05:49 +00:00
|
|
|
not fsub.mesh.hasPeerID(key, receiver.peerInfo.peerId)) and
|
2020-04-21 01:24:42 +00:00
|
|
|
(not fsub.fanout.hasKey(key) or
|
2020-08-12 00:05:49 +00:00
|
|
|
not fsub.fanout.hasPeerID(key , receiver.peerInfo.peerId)):
|
2020-05-27 18:33:49 +00:00
|
|
|
trace "waitSub sleeping..."
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2020-09-21 09:16:29 +00:00
|
|
|
# await more heartbeats
|
|
|
|
await ev.wait()
|
|
|
|
ev.clear()
|
|
|
|
|
2020-04-21 01:24:42 +00:00
|
|
|
dec ceil
|
|
|
|
doAssert(ceil > 0, "waitSub timeout!")
|
|
|
|
|
2020-07-08 00:33:05 +00:00
|
|
|
template tryPublish(call: untyped, require: int, wait: Duration = 1.seconds, times: int = 10): untyped =
|
|
|
|
var
|
|
|
|
limit = times
|
|
|
|
pubs = 0
|
|
|
|
while pubs < require and limit > 0:
|
|
|
|
pubs = pubs + call
|
|
|
|
await sleepAsync(wait)
|
|
|
|
limit.dec()
|
|
|
|
if limit == 0:
|
|
|
|
doAssert(false, "Failed to publish!")
|
|
|
|
|
2019-12-06 02:16:18 +00:00
|
|
|
suite "GossipSub":
|
2020-04-21 01:24:42 +00:00
|
|
|
teardown:
|
2020-09-21 17:48:19 +00:00
|
|
|
checkTrackers()
|
2020-04-21 01:24:42 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
asyncTest "GossipSub validation should succeed":
|
|
|
|
var handlerFut = newFuture[bool]()
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
check topic == "foobar"
|
|
|
|
handlerFut.complete(true)
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
let
|
|
|
|
nodes = generateNodes(2, gossip = true)
|
2020-07-27 19:33:51 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
2020-08-12 00:05:49 +00:00
|
|
|
)
|
2020-07-27 19:33:51 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2020-01-07 08:06:27 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await subscribeNodes(nodes)
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[0].subscribe("foobar", handler)
|
|
|
|
nodes[1].subscribe("foobar", handler)
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var subs: seq[Future[void]]
|
|
|
|
subs &= waitSub(nodes[1], nodes[0], "foobar")
|
|
|
|
subs &= waitSub(nodes[0], nodes[1], "foobar")
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(subs)
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var validatorFut = newFuture[bool]()
|
|
|
|
proc validator(topic: string,
|
|
|
|
message: Message):
|
|
|
|
Future[ValidationResult] {.async.} =
|
|
|
|
check topic == "foobar"
|
|
|
|
validatorFut.complete(true)
|
|
|
|
result = ValidationResult.Accept
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
nodes[1].addValidator("foobar", validator)
|
|
|
|
tryPublish await nodes[0].publish("foobar", "Hello!".toBytes()), 1
|
2020-04-21 01:24:42 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check (await validatorFut) and (await handlerFut)
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
2020-09-21 09:16:29 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
2020-09-21 09:16:29 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
2020-09-21 09:16:29 +00:00
|
|
|
|
2020-12-15 01:25:22 +00:00
|
|
|
asyncTest "GossipSub validation should fail (reject)":
|
2020-11-13 03:44:02 +00:00
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
check false # if we get here, it should fail
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
let
|
|
|
|
nodes = generateNodes(2, gossip = true)
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
2020-08-12 00:05:49 +00:00
|
|
|
)
|
2020-07-27 19:33:51 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await subscribeNodes(nodes)
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[0].subscribe("foobar", handler)
|
|
|
|
nodes[1].subscribe("foobar", handler)
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var subs: seq[Future[void]]
|
|
|
|
subs &= waitSub(nodes[1], nodes[0], "foobar")
|
|
|
|
subs &= waitSub(nodes[0], nodes[1], "foobar")
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(subs)
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
let gossip1 = GossipSub(nodes[0])
|
|
|
|
let gossip2 = GossipSub(nodes[1])
|
2019-12-17 05:24:03 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check:
|
|
|
|
gossip1.mesh["foobar"].len == 1 and "foobar" notin gossip1.fanout
|
|
|
|
gossip2.mesh["foobar"].len == 1 and "foobar" notin gossip2.fanout
|
|
|
|
|
|
|
|
var validatorFut = newFuture[bool]()
|
|
|
|
proc validator(topic: string,
|
|
|
|
message: Message):
|
|
|
|
Future[ValidationResult] {.async.} =
|
|
|
|
result = ValidationResult.Reject
|
|
|
|
validatorFut.complete(true)
|
|
|
|
|
|
|
|
nodes[1].addValidator("foobar", validator)
|
|
|
|
tryPublish await nodes[0].publish("foobar", "Hello!".toBytes()), 1
|
|
|
|
|
|
|
|
check (await validatorFut) == true
|
2020-12-15 01:25:22 +00:00
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
|
|
|
|
asyncTest "GossipSub validation should fail (ignore)":
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
check false # if we get here, it should fail
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(2, gossip = true)
|
|
|
|
|
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
|
|
|
)
|
|
|
|
|
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
|
|
|
|
|
|
|
await subscribeNodes(nodes)
|
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[0].subscribe("foobar", handler)
|
|
|
|
nodes[1].subscribe("foobar", handler)
|
2020-12-15 01:25:22 +00:00
|
|
|
|
|
|
|
var subs: seq[Future[void]]
|
|
|
|
subs &= waitSub(nodes[1], nodes[0], "foobar")
|
|
|
|
subs &= waitSub(nodes[0], nodes[1], "foobar")
|
|
|
|
|
|
|
|
await allFuturesThrowing(subs)
|
|
|
|
|
|
|
|
let gossip1 = GossipSub(nodes[0])
|
|
|
|
let gossip2 = GossipSub(nodes[1])
|
|
|
|
|
|
|
|
check:
|
|
|
|
gossip1.mesh["foobar"].len == 1 and "foobar" notin gossip1.fanout
|
|
|
|
gossip2.mesh["foobar"].len == 1 and "foobar" notin gossip2.fanout
|
|
|
|
|
|
|
|
var validatorFut = newFuture[bool]()
|
|
|
|
proc validator(topic: string,
|
|
|
|
message: Message):
|
|
|
|
Future[ValidationResult] {.async.} =
|
|
|
|
result = ValidationResult.Ignore
|
|
|
|
validatorFut.complete(true)
|
|
|
|
|
|
|
|
nodes[1].addValidator("foobar", validator)
|
|
|
|
tryPublish await nodes[0].publish("foobar", "Hello!".toBytes()), 1
|
|
|
|
|
|
|
|
check (await validatorFut) == true
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
|
|
|
|
asyncTest "GossipSub validation one fails and one succeeds":
|
|
|
|
var handlerFut = newFuture[bool]()
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
check topic == "foo"
|
|
|
|
handlerFut.complete(true)
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(2, gossip = true)
|
|
|
|
|
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
2020-08-12 00:05:49 +00:00
|
|
|
)
|
2020-08-03 05:20:11 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2020-08-03 05:20:11 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await subscribeNodes(nodes)
|
2020-08-03 05:20:11 +00:00
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[1].subscribe("foo", handler)
|
|
|
|
nodes[1].subscribe("bar", handler)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var passed, failed: Future[bool] = newFuture[bool]()
|
|
|
|
proc validator(topic: string,
|
|
|
|
message: Message):
|
|
|
|
Future[ValidationResult] {.async.} =
|
|
|
|
result = if topic == "foo":
|
|
|
|
passed.complete(true)
|
|
|
|
ValidationResult.Accept
|
|
|
|
else:
|
|
|
|
failed.complete(true)
|
|
|
|
ValidationResult.Reject
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
nodes[1].addValidator("foo", "bar", validator)
|
|
|
|
tryPublish await nodes[0].publish("foo", "Hello!".toBytes()), 1
|
|
|
|
tryPublish await nodes[0].publish("bar", "Hello!".toBytes()), 1
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check ((await passed) and (await failed) and (await handlerFut))
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
let gossip1 = GossipSub(nodes[0])
|
|
|
|
let gossip2 = GossipSub(nodes[1])
|
2020-07-27 19:33:51 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check:
|
|
|
|
"foo" notin gossip1.mesh and gossip1.fanout["foo"].len == 1
|
|
|
|
"foo" notin gossip2.mesh and "foo" notin gossip2.fanout
|
|
|
|
"bar" notin gossip1.mesh and gossip1.fanout["bar"].len == 1
|
|
|
|
"bar" notin gossip2.mesh and "bar" notin gossip2.fanout
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
|
|
|
|
asyncTest "e2e - GossipSub should add remote peer topic subscriptions":
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
discard
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(
|
|
|
|
2,
|
2021-02-08 20:33:34 +00:00
|
|
|
gossip = true)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
2020-08-12 00:05:49 +00:00
|
|
|
)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await subscribeNodes(nodes)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[1].subscribe("foobar", handler)
|
2020-11-13 03:44:02 +00:00
|
|
|
await sleepAsync(10.seconds)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
let gossip1 = GossipSub(nodes[0])
|
|
|
|
let gossip2 = GossipSub(nodes[1])
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check:
|
|
|
|
"foobar" in gossip2.topics
|
|
|
|
"foobar" in gossip1.gossipsub
|
|
|
|
gossip1.gossipsub.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
|
|
|
|
asyncTest "e2e - GossipSub should add remote peer topic subscriptions if both peers are subscribed":
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
discard
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(
|
|
|
|
2,
|
2021-02-08 20:33:34 +00:00
|
|
|
gossip = true)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
2020-08-12 00:05:49 +00:00
|
|
|
)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await subscribeNodes(nodes)
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[0].subscribe("foobar", handler)
|
|
|
|
nodes[1].subscribe("foobar", handler)
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var subs: seq[Future[void]]
|
|
|
|
subs &= waitSub(nodes[1], nodes[0], "foobar")
|
|
|
|
subs &= waitSub(nodes[0], nodes[1], "foobar")
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(subs)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
let
|
|
|
|
gossip1 = GossipSub(nodes[0])
|
|
|
|
gossip2 = GossipSub(nodes[1])
|
2020-04-30 13:22:31 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check:
|
|
|
|
"foobar" in gossip1.topics
|
|
|
|
"foobar" in gossip2.topics
|
|
|
|
|
|
|
|
"foobar" in gossip1.gossipsub
|
|
|
|
"foobar" in gossip2.gossipsub
|
|
|
|
|
|
|
|
gossip1.gossipsub.hasPeerID("foobar", gossip2.peerInfo.peerId) or
|
|
|
|
gossip1.mesh.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
|
|
|
|
|
|
|
gossip2.gossipsub.hasPeerID("foobar", gossip1.peerInfo.peerId) or
|
|
|
|
gossip2.mesh.hasPeerID("foobar", gossip1.peerInfo.peerId)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
|
|
|
|
asyncTest "e2e - GossipSub send over fanout A -> B":
|
|
|
|
var passed = newFuture[void]()
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
check topic == "foobar"
|
|
|
|
passed.complete()
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(
|
|
|
|
2,
|
2021-02-08 20:33:34 +00:00
|
|
|
gossip = true)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
|
|
|
)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2020-04-21 01:24:42 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await subscribeNodes(nodes)
|
2020-04-21 01:24:42 +00:00
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[1].subscribe("foobar", handler)
|
2020-11-13 03:44:02 +00:00
|
|
|
await waitSub(nodes[0], nodes[1], "foobar")
|
2020-07-27 19:33:51 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var observed = 0
|
|
|
|
let
|
|
|
|
obs1 = PubSubObserver(onRecv: proc(peer: PubSubPeer; msgs: var RPCMsg) =
|
|
|
|
inc observed
|
2020-08-12 00:05:49 +00:00
|
|
|
)
|
2020-11-13 03:44:02 +00:00
|
|
|
obs2 = PubSubObserver(onSend: proc(peer: PubSubPeer; msgs: var RPCMsg) =
|
|
|
|
inc observed
|
2020-08-12 00:05:49 +00:00
|
|
|
)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
nodes[1].addObserver(obs1)
|
|
|
|
nodes[0].addObserver(obs2)
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
tryPublish await nodes[0].publish("foobar", "Hello!".toBytes()), 1
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var gossip1: GossipSub = GossipSub(nodes[0])
|
|
|
|
var gossip2: GossipSub = GossipSub(nodes[1])
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check:
|
|
|
|
"foobar" in gossip1.gossipsub
|
|
|
|
gossip1.fanout.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
|
|
|
not gossip1.mesh.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
|
|
|
|
|
|
|
await passed.wait(2.seconds)
|
|
|
|
|
|
|
|
trace "test done, stopping..."
|
|
|
|
|
|
|
|
await nodes[0].stop()
|
|
|
|
await nodes[1].stop()
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
check observed == 2
|
|
|
|
|
|
|
|
asyncTest "e2e - GossipSub send over mesh A -> B":
|
|
|
|
var passed: Future[bool] = newFuture[bool]()
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
check topic == "foobar"
|
|
|
|
passed.complete(true)
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(
|
|
|
|
2,
|
2021-02-08 20:33:34 +00:00
|
|
|
gossip = true)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
|
|
|
)
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2020-08-12 00:05:49 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await subscribeNodes(nodes)
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
nodes[0].subscribe("foobar", handler)
|
|
|
|
nodes[1].subscribe("foobar", handler)
|
2020-11-13 03:44:02 +00:00
|
|
|
await waitSub(nodes[0], nodes[1], "foobar")
|
2020-01-07 08:06:27 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
tryPublish await nodes[0].publish("foobar", "Hello!".toBytes()), 1
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check await passed
|
2020-04-21 01:24:42 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
var gossip1: GossipSub = GossipSub(nodes[0])
|
|
|
|
var gossip2: GossipSub = GossipSub(nodes[1])
|
2019-12-06 02:16:18 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check:
|
|
|
|
"foobar" in gossip1.gossipsub
|
|
|
|
"foobar" in gossip2.gossipsub
|
|
|
|
gossip1.mesh.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
|
|
|
not gossip1.fanout.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
|
|
|
gossip2.mesh.hasPeerID("foobar", gossip1.peerInfo.peerId)
|
|
|
|
not gossip2.fanout.hasPeerID("foobar", gossip1.peerInfo.peerId)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
asyncTest "e2e - GossipSub send over floodPublish A -> B":
|
|
|
|
var passed: Future[bool] = newFuture[bool]()
|
|
|
|
proc handler(topic: string, data: seq[byte]) {.async, gcsafe.} =
|
|
|
|
check topic == "foobar"
|
|
|
|
passed.complete(true)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
let
|
2021-04-22 09:51:22 +00:00
|
|
|
nodes = generateNodes(
|
|
|
|
2,
|
|
|
|
gossip = true)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
# start switches
|
|
|
|
nodesFut = await allFinished(
|
|
|
|
nodes[0].switch.start(),
|
|
|
|
nodes[1].switch.start(),
|
|
|
|
)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
# start pubsub
|
|
|
|
await allFuturesThrowing(
|
|
|
|
allFinished(
|
|
|
|
nodes[0].start(),
|
|
|
|
nodes[1].start(),
|
|
|
|
))
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
var gossip1: GossipSub = GossipSub(nodes[0])
|
|
|
|
gossip1.parameters.floodPublish = true
|
|
|
|
var gossip2: GossipSub = GossipSub(nodes[1])
|
|
|
|
gossip2.parameters.floodPublish = true
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
await subscribeNodes(nodes)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
# nodes[0].subscribe("foobar", handler)
|
|
|
|
nodes[1].subscribe("foobar", handler)
|
|
|
|
await waitSub(nodes[0], nodes[1], "foobar")
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
tryPublish await nodes[0].publish("foobar", "Hello!".toBytes()), 1
|
2020-07-20 06:55:00 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
check await passed
|
|
|
|
|
|
|
|
check:
|
|
|
|
"foobar" in gossip1.gossipsub
|
|
|
|
"foobar" notin gossip2.gossipsub
|
|
|
|
not gossip1.mesh.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
|
|
|
not gossip1.fanout.hasPeerID("foobar", gossip2.peerInfo.peerId)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
await allFuturesThrowing(
|
2021-04-22 09:51:22 +00:00
|
|
|
nodes[0].switch.stop(),
|
|
|
|
nodes[1].switch.stop()
|
|
|
|
)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes[0].stop(),
|
|
|
|
nodes[1].stop()
|
|
|
|
)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
await allFuturesThrowing(nodesFut.concat())
|
|
|
|
|
|
|
|
asyncTest "e2e - GossipSub with multiple peers":
|
2020-11-13 03:44:02 +00:00
|
|
|
var runs = 10
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(runs, gossip = true, triggerSelf = true)
|
|
|
|
nodesFut = nodes.mapIt(it.switch.start())
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodes.mapIt(it.start()))
|
2021-04-22 09:51:22 +00:00
|
|
|
await subscribeNodes(nodes)
|
2020-11-13 03:44:02 +00:00
|
|
|
|
|
|
|
var seen: Table[string, int]
|
|
|
|
var seenFut = newFuture[void]()
|
2020-12-14 21:22:53 +00:00
|
|
|
for i in 0..<nodes.len:
|
|
|
|
let dialer = nodes[i]
|
2020-11-13 03:44:02 +00:00
|
|
|
var handler: TopicHandler
|
|
|
|
closureScope:
|
|
|
|
var peerName = $dialer.peerInfo.peerId
|
|
|
|
handler = proc(topic: string, data: seq[byte]) {.async, gcsafe, closure.} =
|
|
|
|
if peerName notin seen:
|
|
|
|
seen[peerName] = 0
|
|
|
|
seen[peerName].inc
|
|
|
|
check topic == "foobar"
|
|
|
|
if not seenFut.finished() and seen.len >= runs:
|
|
|
|
seenFut.complete()
|
|
|
|
|
2020-12-19 14:43:32 +00:00
|
|
|
dialer.subscribe("foobar", handler)
|
2020-11-13 03:44:02 +00:00
|
|
|
await waitSub(nodes[0], dialer, "foobar")
|
|
|
|
|
|
|
|
tryPublish await wait(nodes[0].publish("foobar",
|
|
|
|
toBytes("from node " &
|
|
|
|
$nodes[0].peerInfo.peerId)),
|
|
|
|
1.minutes), 1, 5.seconds
|
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
await wait(seenFut, 2.minutes)
|
2020-11-13 03:44:02 +00:00
|
|
|
check: seen.len >= runs
|
|
|
|
for k, v in seen.pairs:
|
|
|
|
check: v >= 1
|
|
|
|
|
|
|
|
for node in nodes:
|
|
|
|
var gossip = GossipSub(node)
|
2021-04-22 09:51:22 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
check:
|
|
|
|
"foobar" in gossip.gossipsub
|
2020-07-27 19:33:51 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes.mapIt(
|
|
|
|
allFutures(
|
|
|
|
it.stop(),
|
|
|
|
it.switch.stop())))
|
2020-01-10 03:59:27 +00:00
|
|
|
|
2020-11-13 03:44:02 +00:00
|
|
|
await allFuturesThrowing(nodesFut)
|
2021-01-13 14:49:44 +00:00
|
|
|
|
2021-04-22 09:51:22 +00:00
|
|
|
asyncTest "e2e - GossipSub with multiple peers (sparse)":
|
2021-02-26 05:15:58 +00:00
|
|
|
var runs = 10
|
|
|
|
|
|
|
|
let
|
|
|
|
nodes = generateNodes(runs, gossip = true, triggerSelf = true)
|
|
|
|
nodesFut = nodes.mapIt(it.switch.start())
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodes.mapIt(it.start()))
|
2021-04-22 09:51:22 +00:00
|
|
|
await subscribeSparseNodes(nodes)
|
2021-02-26 05:15:58 +00:00
|
|
|
|
|
|
|
var seen: Table[string, int]
|
|
|
|
var seenFut = newFuture[void]()
|
|
|
|
for i in 0..<nodes.len:
|
|
|
|
let dialer = nodes[i]
|
|
|
|
var handler: TopicHandler
|
|
|
|
closureScope:
|
|
|
|
var peerName = $dialer.peerInfo.peerId
|
|
|
|
handler = proc(topic: string, data: seq[byte]) {.async, gcsafe, closure.} =
|
|
|
|
if peerName notin seen:
|
|
|
|
seen[peerName] = 0
|
|
|
|
seen[peerName].inc
|
|
|
|
check topic == "foobar"
|
|
|
|
if not seenFut.finished() and seen.len >= runs:
|
|
|
|
seenFut.complete()
|
|
|
|
|
|
|
|
dialer.subscribe("foobar", handler)
|
|
|
|
await waitSub(nodes[0], dialer, "foobar")
|
|
|
|
|
|
|
|
tryPublish await wait(nodes[0].publish("foobar",
|
|
|
|
toBytes("from node " &
|
|
|
|
$nodes[0].peerInfo.peerId)),
|
|
|
|
1.minutes), 1, 5.seconds
|
|
|
|
|
|
|
|
await wait(seenFut, 5.minutes)
|
|
|
|
check: seen.len >= runs
|
|
|
|
for k, v in seen.pairs:
|
|
|
|
check: v >= 1
|
|
|
|
|
|
|
|
for node in nodes:
|
|
|
|
var gossip = GossipSub(node)
|
|
|
|
check:
|
|
|
|
"foobar" in gossip.gossipsub
|
|
|
|
gossip.fanout.len == 0
|
|
|
|
gossip.mesh["foobar"].len > 0
|
|
|
|
|
|
|
|
await allFuturesThrowing(
|
|
|
|
nodes.mapIt(
|
|
|
|
allFutures(
|
|
|
|
it.stop(),
|
|
|
|
it.switch.stop())))
|
|
|
|
|
|
|
|
await allFuturesThrowing(nodesFut)
|