mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-26 14:33:27 +00:00
feat(mix): DoS protection + libp2p v2.0.0 + stateless RLN + tests (rebased onto #3935)
Squash of 13 commits from feat/mix-dos-protection-libp2p-v2.0.0 onto the logos_delivery/ folder-restructure base from #3935 (build-messaging-folder). Original commit history (squashed): - d8e6dcef feat(mix): integrate mix protocol with extended kademlia + RLN spam protection - fb72f18d refactor(mix): split DoS-protection self-registration into background retry - d8bbef0c feat(mix): bump libp2p stack to v2.0.0 + adopt stateless RLN spam protection - 2f24448a fix(tests): use HmacDrbgContext.new() instead of crypto.newRng() - 5a21455c fix(ci): regen nimble.lock for v2.0.0 + disambiguate rng in wakucore - 03ef02a2 fix(tests): wrap HmacDrbgContext via newBearSslRng for libp2p v2.0.0 - 167ab1df fix(nix): regenerate deps.nix from updated nimble.lock - 97a27222 fix(tests): wrap or pass Rng correctly for 3-arg PrivateKey.random - 5561fcb5 fix(tests): replace removed newStandardSwitch with SwitchBuilder - ba39ee4a fix(tests): libp2p v2.0.0 API migrations across test suite - 328e11df fix: gitignore test binaries + remove accidentally-committed binary - cc712444 fix(tests): more v2.0.0 API migrations (rng template, PeerId.random, etc.) - 412d97a9 fix(tests): unblock CI — nph, excise orphan waku_noise, complete v2.0.0 Rng migration Conflict resolutions (#3935 → ours): - 11 import-path migrations: waku/X → logos_delivery/waku/X - waku_node/waku_node/relay.nim: dropped our `registerRelayHandler` proc (relocated to subscription_manager.nim by #3935; see cascade fix below) - factory/builder.nim: combined both sides' new imports (net_config + waku_switch) - factory/conf_builder/mix_conf_builder.nim: libp2p_mix package (not libp2p/protocols/mix) - waku_mix/protocol.nim: combined paths + our mix_rln_spam_protection/relay/nimchronos imports - 3 test files: dropped noise_utils import (replicates noise excision from original PR) - 2 UA file moves: option_shims.nim and waku_mix_coordination.nim added at new paths Cascade fixes (#3935 lost our work, restored): - subscription_manager.nim: added `mixHandler` to #3935's `registerRelayHandler`, and added `waku_mix` to its imports. Without this, mix messages were silently dropped from the relay handler chain. - config.nims: option_shims auto-import path migrated to logos_delivery/... Validation: - nph check on all 82 staged .nim files: clean (0 reformats needed) - wakunode2 build: exit 0, 38 MB binary - (sim PASS confirmed in earlier identical-state run: 5/5 mix init, 5 RLN proofs gen/verify, 0 errors) Backup tag at original tip: backup/3931-pre-3935-rebase (412d97a9).
This commit is contained in:
parent
2fe7e1c373
commit
bd65c116bb
9
.gitignore
vendored
9
.gitignore
vendored
@ -3,6 +3,11 @@
|
||||
# Executables shall be put in an ignored build/ directory
|
||||
/build
|
||||
|
||||
# Test binaries (built by `nim c tests/...nim` for local debug compile)
|
||||
/tests/all_tests_common
|
||||
/tests/all_tests_waku
|
||||
/tests/all_tests_wakunode2
|
||||
|
||||
# Generated Files
|
||||
*.generated.nim
|
||||
|
||||
@ -41,6 +46,7 @@ node_modules/
|
||||
|
||||
# RLN / keystore
|
||||
rlnKeystore.json
|
||||
rln_keystore*.json
|
||||
*.tar.gz
|
||||
|
||||
# sqlite db
|
||||
@ -91,3 +97,6 @@ nimbledeps
|
||||
# Python bytecode from tests/simulator
|
||||
__pycache__/
|
||||
*.pyc
|
||||
|
||||
# sim driver script (local dev tool, not part of build/CI)
|
||||
simulations/mixnet/roundtrip_check.sh
|
||||
|
||||
@ -31,9 +31,8 @@ import
|
||||
protocols/kademlia/types,
|
||||
protocols/service_discovery/types as sd_types,
|
||||
nameresolving/dnsresolver,
|
||||
protocols/mix/curve25519,
|
||||
protocols/mix/mix_protocol,
|
||||
] # define DNS resolution
|
||||
import libp2p_mix/curve25519, libp2p_mix/mix_protocol
|
||||
import
|
||||
logos_delivery/waku/[
|
||||
waku_core,
|
||||
@ -51,6 +50,8 @@ import
|
||||
common/utils/nat,
|
||||
waku_store/common,
|
||||
waku_filter_v2/client,
|
||||
waku_filter_v2/common as filter_common,
|
||||
waku_mix/protocol,
|
||||
common/logging,
|
||||
],
|
||||
./config_chat2mix
|
||||
@ -61,6 +62,99 @@ import ../../logos_delivery/waku/waku_rln_relay
|
||||
logScope:
|
||||
topics = "chat2 mix"
|
||||
|
||||
#########################
|
||||
## Mix Spam Protection ##
|
||||
#########################
|
||||
|
||||
# Forward declaration
|
||||
proc maintainSpamProtectionSubscription(
|
||||
node: WakuNode, contentTopics: seq[ContentTopic]
|
||||
) {.async.}
|
||||
|
||||
proc setupMixSpamProtectionViaFilter(node: WakuNode) {.async.} =
|
||||
# Register message handler for spam protection coordination
|
||||
let spamTopics = node.wakuMix.getSpamProtectionContentTopics()
|
||||
|
||||
proc handleSpamMessage(
|
||||
pubsubTopic: PubsubTopic, message: WakuMessage
|
||||
): Future[void] {.async, gcsafe.} =
|
||||
await node.wakuMix.handleMessage(pubsubTopic, message)
|
||||
|
||||
node.wakuFilterClient.registerPushHandler(handleSpamMessage)
|
||||
|
||||
# Wait for filter peer and maintain subscription
|
||||
asyncSpawn maintainSpamProtectionSubscription(node, spamTopics)
|
||||
|
||||
proc maintainSpamProtectionSubscription(
|
||||
node: WakuNode, contentTopics: seq[ContentTopic]
|
||||
) {.async.} =
|
||||
const RetryInterval = chronos.seconds(5)
|
||||
const SubscriptionMaintenance = chronos.seconds(30)
|
||||
const MaxFailedSubscribes = 3
|
||||
var currentFilterPeer: Option[RemotePeerInfo] = none(RemotePeerInfo)
|
||||
var noFailedSubscribes = 0
|
||||
|
||||
while true:
|
||||
# Select or reuse filter peer
|
||||
if currentFilterPeer.isNone():
|
||||
let filterPeerOpt = node.peerManager.selectPeer(WakuFilterSubscribeCodec)
|
||||
if filterPeerOpt.isNone():
|
||||
debug "No filter peer available yet for spam protection, retrying..."
|
||||
await sleepAsync(RetryInterval)
|
||||
continue
|
||||
currentFilterPeer = some(filterPeerOpt.get())
|
||||
info "Selected filter peer for spam protection",
|
||||
peer = currentFilterPeer.get().peerId
|
||||
|
||||
# Check if subscription is still alive with ping
|
||||
let pingErr = (await node.wakuFilterClient.ping(currentFilterPeer.get())).errorOr:
|
||||
# Subscription is alive, wait before next check
|
||||
await sleepAsync(SubscriptionMaintenance)
|
||||
if noFailedSubscribes > 0:
|
||||
noFailedSubscribes = 0
|
||||
continue
|
||||
|
||||
# Subscription lost, need to re-subscribe
|
||||
warn "Spam protection filter subscription ping failed, re-subscribing",
|
||||
error = pingErr, peer = currentFilterPeer.get().peerId
|
||||
|
||||
# Determine pubsub topic from content topics (using auto-sharding)
|
||||
if node.wakuAutoSharding.isNone():
|
||||
error "Auto-sharding not configured, cannot determine pubsub topic for spam protection"
|
||||
await sleepAsync(RetryInterval)
|
||||
continue
|
||||
|
||||
let shardRes = node.wakuAutoSharding.get().getShard(contentTopics[0])
|
||||
if shardRes.isErr():
|
||||
error "Failed to determine shard for spam protection", error = shardRes.error
|
||||
await sleepAsync(RetryInterval)
|
||||
continue
|
||||
|
||||
let shard = shardRes.get()
|
||||
let pubsubTopic: PubsubTopic = shard # converter toPubsubTopic
|
||||
|
||||
# Subscribe to spam protection topics
|
||||
let res = await node.wakuFilterClient.subscribe(
|
||||
currentFilterPeer.get(), pubsubTopic, contentTopics
|
||||
)
|
||||
if res.isErr():
|
||||
noFailedSubscribes += 1
|
||||
warn "Failed to subscribe to spam protection topics via filter",
|
||||
error = res.error, topics = contentTopics, failCount = noFailedSubscribes
|
||||
|
||||
if noFailedSubscribes >= MaxFailedSubscribes:
|
||||
# Try with a different peer
|
||||
warn "Max subscription failures reached, selecting new filter peer"
|
||||
currentFilterPeer = none(RemotePeerInfo)
|
||||
noFailedSubscribes = 0
|
||||
|
||||
await sleepAsync(RetryInterval)
|
||||
else:
|
||||
info "Successfully subscribed to spam protection topics via filter",
|
||||
topics = contentTopics, peer = currentFilterPeer.get().peerId
|
||||
noFailedSubscribes = 0
|
||||
await sleepAsync(SubscriptionMaintenance)
|
||||
|
||||
const Help = """
|
||||
Commands: /[?|help|connect|nick|exit]
|
||||
help: Prints this help
|
||||
@ -213,20 +307,21 @@ proc publish(c: Chat, line: string) {.async.} =
|
||||
try:
|
||||
if not c.node.wakuLightpushClient.isNil():
|
||||
# Attempt lightpush with mix
|
||||
|
||||
(
|
||||
waitFor c.node.lightpushPublish(
|
||||
some(c.conf.getPubsubTopic(c.node, c.contentTopic)),
|
||||
message,
|
||||
none(RemotePeerInfo),
|
||||
true,
|
||||
)
|
||||
).isOkOr:
|
||||
error "failed to publish lightpush message", error = error
|
||||
let res = await c.node.lightpushPublish(
|
||||
some(c.conf.getPubsubTopic(c.node, c.contentTopic)),
|
||||
message,
|
||||
none(RemotePeerInfo),
|
||||
true,
|
||||
)
|
||||
if res.isErr():
|
||||
error "failed to publish lightpush message", error = res.error
|
||||
echo "Error: " & res.error.desc.get("unknown error")
|
||||
else:
|
||||
error "failed to publish message as lightpush client is not initialized"
|
||||
echo "Error: lightpush client is not initialized"
|
||||
except CatchableError:
|
||||
error "caught error publishing message: ", error = getCurrentExceptionMsg()
|
||||
echo "Error: " & getCurrentExceptionMsg()
|
||||
|
||||
# TODO This should read or be subscribe handler subscribe
|
||||
proc readAndPrint(c: Chat) {.async.} =
|
||||
@ -455,7 +550,11 @@ proc processInput(rfd: AsyncFD, rng: crypto.Rng) {.async.} =
|
||||
error "failed to generate mix key pair", error = error
|
||||
return
|
||||
|
||||
(await node.mountMix(conf.clusterId, mixPrivKey, conf.mixnodes)).isOkOr:
|
||||
(
|
||||
await node.mountMix(
|
||||
conf.clusterId, mixPrivKey, conf.mixnodes, some(conf.rlnUserMessageLimit)
|
||||
)
|
||||
).isOkOr:
|
||||
error "failed to mount waku mix protocol: ", error = $error
|
||||
quit(QuitFailure)
|
||||
|
||||
@ -488,6 +587,10 @@ proc processInput(rfd: AsyncFD, rng: crypto.Rng) {.async.} =
|
||||
|
||||
#await node.mountRendezvousClient(conf.clusterId)
|
||||
|
||||
# Subscribe to spam protection coordination topics via filter since chat2mix doesn't use relay
|
||||
if not node.wakuFilterClient.isNil():
|
||||
asyncSpawn setupMixSpamProtectionViaFilter(node)
|
||||
|
||||
await node.start()
|
||||
|
||||
node.peerManager.start()
|
||||
|
||||
@ -240,6 +240,13 @@ type
|
||||
name: "kad-bootstrap-node"
|
||||
.}: seq[string]
|
||||
|
||||
## RLN spam protection config
|
||||
rlnUserMessageLimit* {.
|
||||
desc: "Maximum messages per epoch for RLN spam protection.",
|
||||
defaultValue: 100,
|
||||
name: "rln-user-message-limit"
|
||||
.}: int
|
||||
|
||||
proc parseCmdArg*(T: type MixNodePubInfo, p: string): T =
|
||||
let elements = p.split(":")
|
||||
if elements.len != 2:
|
||||
|
||||
@ -111,6 +111,8 @@ if not defined(macosx) and not defined(android):
|
||||
nimStackTraceOverride
|
||||
switch("import", "libbacktrace")
|
||||
|
||||
switch("import", "logos_delivery/waku/compat/option_valueor")
|
||||
|
||||
--define:
|
||||
nimOldCaseObjects
|
||||
# https://github.com/status-im/nim-confutils/issues/9
|
||||
|
||||
@ -62,6 +62,15 @@ requires "nim >= 2.2.4",
|
||||
# Packages not on nimble (use git URLs)
|
||||
|
||||
requires "https://github.com/logos-messaging/nim-ffi#v0.1.3"
|
||||
requires "https://github.com/logos-co/mix-rln-spam-protection-plugin.git#23b278b4ab21193ad4e9ce76015f008db7332a6f"
|
||||
|
||||
# nim-libp2p-mix: extracted mix protocol used by the plugin and by waku's
|
||||
# mix integration layer. Tip of experiment/drop-nimble-lock (PR #14, stacked
|
||||
# on chore/bump-libp2p-v2.0.0). Carries the v2.0.0 bump + sink overrides +
|
||||
# AddressConfidence.Infinite + deeper move-semantics propagation + the
|
||||
# lockfile-as-build-artefact cleanup. Re-bump to master SHA once #14 lands.
|
||||
# The plugin pins the same SHA — keeps the diamond dep collapsed.
|
||||
requires "https://github.com/logos-co/nim-libp2p-mix.git#50c4ab4fa788a33eb12a0a2cecaa708873352b58"
|
||||
|
||||
requires "https://github.com/logos-messaging/nim-sds.git#b12f5ee07c5b764303b51fb948b32a4ade1de3b5"
|
||||
|
||||
@ -69,7 +78,6 @@ requires "https://github.com/NagyZoltanPeter/nim-brokers.git#v3.1.1"
|
||||
|
||||
requires "https://github.com/vacp2p/nim-lsquic.git#v0.5.1"
|
||||
requires "https://github.com/vacp2p/nim-jwt.git#057ec95eb5af0eea9c49bfe9025b3312c95dc5f2"
|
||||
requires "https://github.com/logos-co/nim-libp2p-mix#380513117d556bf8f70066f5e72a7fd74fe36ba6"
|
||||
|
||||
proc getMyCPU(): string =
|
||||
## Need to set cpu more explicit manner to avoid arch issues between dependencies
|
||||
|
||||
@ -111,8 +111,21 @@ proc setupSwitchServices(
|
||||
MaxNumRelayServers, RelayClient(circuitRelay), onReservation, rng
|
||||
)
|
||||
let holePunchService = HPService.new(autonatService, autoRelayService)
|
||||
# libp2p v2.0.0: switch.start() no longer auto-calls service.setup() (part
|
||||
# of the Service lifecycle refactor in libp2p#2462). Without setup,
|
||||
# HPService's wrapped Autonat/AutoRelay leave their addressMapper field
|
||||
# nil, which makes peerInfo.expandAddrs SIGSEGV during start().
|
||||
try:
|
||||
holePunchService.setup(waku.node.switch)
|
||||
except ServiceSetupError as e:
|
||||
error "HPService setup failed", description = e.msg
|
||||
waku.node.switch.services = @[Service(holePunchService)]
|
||||
else:
|
||||
# Same reason as above: AutonatService.setup() initializes addressMapper.
|
||||
try:
|
||||
autonatService.setup(waku.node.switch)
|
||||
except ServiceSetupError as e:
|
||||
error "AutonatService setup failed", description = e.msg
|
||||
waku.node.switch.services = @[Service(autonatService)]
|
||||
|
||||
# libp2p 2.0.0 split Service.setup out of Service.start: the switch runs setup
|
||||
|
||||
@ -28,6 +28,7 @@ import
|
||||
node/health_monitor/online_monitor,
|
||||
node/waku_switch,
|
||||
],
|
||||
../waku_switch,
|
||||
./peer_store/peer_storage,
|
||||
./waku_peer_store
|
||||
|
||||
|
||||
@ -11,6 +11,7 @@ import
|
||||
node/node_telemetry,
|
||||
waku_relay,
|
||||
waku_archive,
|
||||
waku_mix,
|
||||
waku_store_sync,
|
||||
waku_filter_v2/common as filter_common,
|
||||
waku_filter_v2/client as filter_client,
|
||||
@ -66,6 +67,12 @@ proc registerRelayHandler(
|
||||
|
||||
node.wakuStoreReconciliation.messageIngress(topic, msg)
|
||||
|
||||
proc mixHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
if node.wakuMix.isNil():
|
||||
return
|
||||
|
||||
await node.wakuMix.handleMessage(topic, msg)
|
||||
|
||||
proc internalHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
MessageSeenEvent.emit(node.brokerCtx, topic, msg)
|
||||
|
||||
@ -76,6 +83,7 @@ proc registerRelayHandler(
|
||||
await filterHandler(topic, msg)
|
||||
await archiveHandler(topic, msg)
|
||||
await syncHandler(topic, msg)
|
||||
await mixHandler(topic, msg)
|
||||
await internalHandler(topic, msg)
|
||||
|
||||
if node.legacyAppHandlers.hasKey(topic) and not node.legacyAppHandlers[topic].isNil():
|
||||
|
||||
125
logos_delivery/waku/node/waku_mix_coordination.nim
Normal file
125
logos_delivery/waku/node/waku_mix_coordination.nim
Normal file
@ -0,0 +1,125 @@
|
||||
## Mix spam protection coordination via filter protocol
|
||||
## This module handles filter-based subscription for spam protection coordination
|
||||
## when relay is not available.
|
||||
|
||||
{.push raises: [].}
|
||||
|
||||
import chronos, chronicles, std/options
|
||||
import
|
||||
../waku_core,
|
||||
../waku_core/topics/sharding,
|
||||
../waku_filter_v2/common,
|
||||
./peer_manager,
|
||||
../waku_filter_v2/client,
|
||||
../waku_mix/protocol
|
||||
|
||||
logScope:
|
||||
topics = "waku node mix_coordination"
|
||||
|
||||
# Type aliases for callbacks to avoid circular imports
|
||||
type
|
||||
FilterSubscribeProc* = proc(
|
||||
pubsubTopic: Option[PubsubTopic],
|
||||
contentTopics: seq[ContentTopic],
|
||||
peer: RemotePeerInfo,
|
||||
): Future[FilterSubscribeResult] {.async, gcsafe.}
|
||||
|
||||
FilterPingProc* =
|
||||
proc(peer: RemotePeerInfo): Future[FilterSubscribeResult] {.async, gcsafe.}
|
||||
|
||||
# Forward declaration
|
||||
proc subscribeSpamProtectionViaFilter(
|
||||
wakuMix: WakuMix,
|
||||
peerManager: PeerManager,
|
||||
filterClient: WakuFilterClient,
|
||||
filterSubscribe: FilterSubscribeProc,
|
||||
contentTopics: seq[ContentTopic],
|
||||
) {.async.}
|
||||
|
||||
proc setupSpamProtectionViaFilter*(
|
||||
wakuMix: WakuMix,
|
||||
peerManager: PeerManager,
|
||||
filterClient: WakuFilterClient,
|
||||
filterSubscribe: FilterSubscribeProc,
|
||||
) =
|
||||
## Set up filter-based spam protection coordination.
|
||||
## Registers message handler and spawns subscription maintenance task.
|
||||
let spamTopics = wakuMix.getSpamProtectionContentTopics()
|
||||
if spamTopics.len == 0:
|
||||
return
|
||||
|
||||
info "Relay not available, subscribing to spam protection via filter",
|
||||
topics = spamTopics
|
||||
|
||||
# Register handler for spam protection messages
|
||||
filterClient.registerPushHandler(
|
||||
proc(pubsubTopic: PubsubTopic, message: WakuMessage) {.async, gcsafe.} =
|
||||
if message.contentTopic in spamTopics:
|
||||
await wakuMix.handleMessage(pubsubTopic, message)
|
||||
)
|
||||
|
||||
# Wait for filter peer to be available and maintain subscription
|
||||
asyncSpawn subscribeSpamProtectionViaFilter(
|
||||
wakuMix, peerManager, filterClient, filterSubscribe, spamTopics
|
||||
)
|
||||
|
||||
proc subscribeSpamProtectionViaFilter(
|
||||
wakuMix: WakuMix,
|
||||
peerManager: PeerManager,
|
||||
filterClient: WakuFilterClient,
|
||||
filterSubscribe: FilterSubscribeProc,
|
||||
contentTopics: seq[ContentTopic],
|
||||
) {.async.} =
|
||||
## Subscribe to spam protection topics via filter and maintain the subscription.
|
||||
## Waits for a filter peer to be available before subscribing.
|
||||
## Continuously monitors the subscription health with periodic pings.
|
||||
const RetryInterval = chronos.seconds(5)
|
||||
const SubscriptionMaintenance = chronos.seconds(30)
|
||||
const MaxFailedSubscribes = 3
|
||||
var currentFilterPeer: Option[RemotePeerInfo] = none(RemotePeerInfo)
|
||||
var noFailedSubscribes = 0
|
||||
|
||||
while true:
|
||||
# Select or reuse filter peer
|
||||
if currentFilterPeer.isNone():
|
||||
let filterPeerOpt = peerManager.selectPeer(WakuFilterSubscribeCodec)
|
||||
if filterPeerOpt.isNone():
|
||||
debug "No filter peer available yet for spam protection, retrying..."
|
||||
await sleepAsync(RetryInterval)
|
||||
continue
|
||||
currentFilterPeer = some(filterPeerOpt.get())
|
||||
info "Selected filter peer for spam protection",
|
||||
peer = currentFilterPeer.get().peerId
|
||||
|
||||
# Check if subscription is still alive with ping
|
||||
let pingErr = (await filterClient.ping(currentFilterPeer.get())).errorOr:
|
||||
# Subscription is alive, wait before next check
|
||||
await sleepAsync(SubscriptionMaintenance)
|
||||
if noFailedSubscribes > 0:
|
||||
noFailedSubscribes = 0
|
||||
continue
|
||||
|
||||
# Subscription lost, need to re-subscribe
|
||||
warn "Spam protection filter subscription ping failed, re-subscribing",
|
||||
error = pingErr, peer = currentFilterPeer.get().peerId
|
||||
|
||||
# Subscribe to spam protection topics
|
||||
let res =
|
||||
await filterSubscribe(none(PubsubTopic), contentTopics, currentFilterPeer.get())
|
||||
if res.isErr():
|
||||
noFailedSubscribes += 1
|
||||
warn "Failed to subscribe to spam protection topics via filter",
|
||||
error = res.error, topics = contentTopics, failCount = noFailedSubscribes
|
||||
|
||||
if noFailedSubscribes >= MaxFailedSubscribes:
|
||||
# Try with a different peer
|
||||
warn "Max subscription failures reached, selecting new filter peer"
|
||||
currentFilterPeer = none(RemotePeerInfo)
|
||||
noFailedSubscribes = 0
|
||||
|
||||
await sleepAsync(RetryInterval)
|
||||
else:
|
||||
info "Successfully subscribed to spam protection topics via filter",
|
||||
topics = contentTopics, peer = currentFilterPeer.get().peerId
|
||||
noFailedSubscribes = 0
|
||||
await sleepAsync(SubscriptionMaintenance)
|
||||
@ -190,6 +190,18 @@ proc getShardsGetter(node: WakuNode, configuredShards: seq[uint16]): GetShards =
|
||||
return shards
|
||||
return configuredShards
|
||||
|
||||
proc getRelayMixHandler*(node: WakuNode): Option[WakuRelayHandler] =
|
||||
## Returns a handler for mix spam protection coordination messages if mix is mounted
|
||||
if node.wakuMix.isNil():
|
||||
return none(WakuRelayHandler)
|
||||
|
||||
let handler: WakuRelayHandler = proc(
|
||||
pubsubTopic: PubsubTopic, message: WakuMessage
|
||||
): Future[void] {.async, gcsafe.} =
|
||||
await node.wakuMix.handleMessage(pubsubTopic, message)
|
||||
|
||||
return some(handler)
|
||||
|
||||
proc getCapabilitiesGetter(node: WakuNode): GetCapabilities =
|
||||
return proc(): seq[Capabilities] {.closure, gcsafe, raises: [].} =
|
||||
if node.wakuRelay.isNil():
|
||||
@ -313,6 +325,7 @@ proc mountMix*(
|
||||
clusterId: uint16,
|
||||
mixPrivKey: Curve25519Key,
|
||||
mixnodes: seq[MixNodePubInfo],
|
||||
userMessageLimit: Option[int] = none(int),
|
||||
): Future[Result[void, string]] {.async.} =
|
||||
info "mounting mix protocol", nodeId = node.info #TODO log the config used
|
||||
|
||||
@ -323,8 +336,29 @@ proc mountMix*(
|
||||
return err("Failed to convert multiaddress to string.")
|
||||
info "local addr", localaddr = localaddrStr
|
||||
|
||||
# Create callback to publish coordination messages via relay
|
||||
let publishMessage: PublishMessage = proc(
|
||||
message: WakuMessage
|
||||
): Future[Result[void, string]] {.async.} =
|
||||
# Inline implementation of publish logic to avoid circular import
|
||||
if node.wakuRelay.isNil():
|
||||
return err("WakuRelay not mounted")
|
||||
|
||||
# Derive pubsub topic from content topic using auto sharding
|
||||
let pubsubTopic =
|
||||
if node.wakuAutoSharding.isNone():
|
||||
return err("Auto sharding not configured")
|
||||
else:
|
||||
node.wakuAutoSharding.get().getShard(message.contentTopic).valueOr:
|
||||
return err("Autosharding error: " & error)
|
||||
|
||||
# Publish via relay
|
||||
discard await node.wakuRelay.publish(pubsubTopic, message)
|
||||
return ok()
|
||||
|
||||
node.wakuMix = WakuMix.new(
|
||||
localaddrStr, node.peerManager, clusterId, mixPrivKey, mixnodes
|
||||
localaddrStr, node.peerManager, clusterId, mixPrivKey, mixnodes, publishMessage,
|
||||
userMessageLimit,
|
||||
).valueOr:
|
||||
error "Waku Mix protocol initialization failed", err = error
|
||||
return
|
||||
@ -334,6 +368,7 @@ proc mountMix*(
|
||||
node.switch.mount(node.wakuMix)
|
||||
catchRes.isOkOr:
|
||||
return err(error.msg)
|
||||
|
||||
return ok()
|
||||
|
||||
proc mountKademlia*(
|
||||
@ -620,9 +655,18 @@ proc start*(node: WakuNode) {.async.} =
|
||||
## NOTE: This will dispatch gossipsub start to the WakuRelay.start method override
|
||||
await node.switch.start()
|
||||
|
||||
if not node.wakuMix.isNil():
|
||||
await node.wakuMix.start()
|
||||
|
||||
# After switch.start, run custom Logos Delivery relay start logic
|
||||
await node.reconnectRelayPeers()
|
||||
|
||||
# Kick off the DoS-protection registration broadcast now that peers are
|
||||
# reconnected. Fire-and-forget: the proc returns immediately and an
|
||||
# internal background task retries until the broadcast lands.
|
||||
if not node.wakuMix.isNil():
|
||||
node.wakuMix.registerDoSProtectionWithNetwork()
|
||||
|
||||
node.started = true
|
||||
|
||||
if not node.wakuKademlia.isNil():
|
||||
|
||||
@ -29,6 +29,7 @@ import
|
||||
waku_archive,
|
||||
waku_store_sync,
|
||||
waku_rln_relay,
|
||||
waku_mix,
|
||||
node/waku_node,
|
||||
node/subscription_manager,
|
||||
node/peer_manager,
|
||||
|
||||
@ -1,7 +1,7 @@
|
||||
import logos_delivery/waku/compat/option_valueor
|
||||
{.push raises: [].}
|
||||
|
||||
import chronicles, std/options, chronos, results, metrics
|
||||
import chronicles, std/[options, sequtils], chronos, results, metrics
|
||||
|
||||
import
|
||||
libp2p/crypto/curve25519,
|
||||
@ -11,14 +11,18 @@ import
|
||||
libp2p_mix/mix_protocol,
|
||||
libp2p_mix/mix_metrics,
|
||||
libp2p_mix/delay_strategy,
|
||||
libp2p/[multiaddress, peerid],
|
||||
libp2p_mix/spam_protection,
|
||||
libp2p/[multiaddress, multicodec, peerid, peerinfo],
|
||||
eth/common/keys
|
||||
|
||||
import
|
||||
logos_delivery/waku/node/peer_manager,
|
||||
logos_delivery/waku/waku_core,
|
||||
logos_delivery/waku/waku_enr,
|
||||
logos_delivery/waku/node/peer_manager/waku_peer_store
|
||||
logos_delivery/waku/node/peer_manager/waku_peer_store,
|
||||
mix_rln_spam_protection,
|
||||
logos_delivery/waku/waku_relay,
|
||||
logos_delivery/waku/common/nimchronos
|
||||
|
||||
logScope:
|
||||
topics = "waku mix"
|
||||
@ -26,10 +30,20 @@ logScope:
|
||||
const minMixPoolSize = 4
|
||||
|
||||
type
|
||||
PublishMessage* = proc(message: WakuMessage): Future[Result[void, string]] {.
|
||||
async, gcsafe, raises: []
|
||||
.}
|
||||
|
||||
WakuMix* = ref object of MixProtocol
|
||||
peerManager*: PeerManager
|
||||
clusterId: uint16
|
||||
pubKey*: Curve25519Key
|
||||
mixRlnSpamProtection*: MixRlnSpamProtection
|
||||
publishMessage*: PublishMessage
|
||||
dosRegistrationTask: Future[void]
|
||||
## Background task that retries DoS-protection self-registration until
|
||||
## it succeeds. nil until kicked off via registerDoSProtectionWithNetwork;
|
||||
## cancelled in stop().
|
||||
|
||||
WakuMixResult*[T] = Result[T, string]
|
||||
|
||||
@ -42,11 +56,9 @@ proc processBootNodes(
|
||||
) =
|
||||
var count = 0
|
||||
for node in bootnodes:
|
||||
let pInfo = parsePeerInfo(node.multiAddr).valueOr:
|
||||
error "Failed to get peer id from multiaddress: ",
|
||||
error = error, multiAddr = $node.multiAddr
|
||||
let (peerId, networkAddr) = parseFullAddress(node.multiAddr).valueOr:
|
||||
error "Failed to parse multiaddress", multiAddr = node.multiAddr, error = error
|
||||
continue
|
||||
let peerId = pInfo.peerId
|
||||
var peerPubKey: crypto.PublicKey
|
||||
if not peerId.extractPublicKey(peerPubKey):
|
||||
warn "Failed to extract public key from peerId, skipping node", peerId = peerId
|
||||
@ -66,10 +78,10 @@ proc processBootNodes(
|
||||
count.inc()
|
||||
|
||||
peermgr.addPeer(
|
||||
RemotePeerInfo.init(peerId, @[multiAddr], mixPubKey = some(node.pubKey))
|
||||
RemotePeerInfo.init(peerId, @[networkAddr], mixPubKey = some(node.pubKey))
|
||||
)
|
||||
mix_pool_size.set(count)
|
||||
info "using mix bootstrap nodes ", count = count
|
||||
debug "using mix bootstrap nodes ", count = count
|
||||
|
||||
proc new*(
|
||||
T: typedesc[WakuMix],
|
||||
@ -78,9 +90,11 @@ proc new*(
|
||||
clusterId: uint16,
|
||||
mixPrivKey: Curve25519Key,
|
||||
bootnodes: seq[MixNodePubInfo],
|
||||
publishMessage: PublishMessage,
|
||||
userMessageLimit: Option[int] = none(int),
|
||||
): WakuMixResult[T] =
|
||||
let mixPubKey = public(mixPrivKey)
|
||||
info "mixPubKey", mixPubKey = mixPubKey
|
||||
trace "mixPubKey", mixPubKey = mixPubKey
|
||||
let nodeMultiAddr = MultiAddress.init(nodeAddr).valueOr:
|
||||
return err("failed to parse mix node address: " & $nodeAddr & ", error: " & error)
|
||||
let localMixNodeInfo = initMixNodeInfo(
|
||||
@ -88,13 +102,34 @@ proc new*(
|
||||
peermgr.switch.peerInfo.publicKey.skkey, peermgr.switch.peerInfo.privateKey.skkey,
|
||||
)
|
||||
|
||||
var m = WakuMix(peerManager: peermgr, clusterId: clusterId, pubKey: mixPubKey)
|
||||
# Initialize spam protection with persistent credentials
|
||||
# Use peerID in keystore path so multiple peers can run from same directory
|
||||
# Tree path is shared across all nodes to maintain the full membership set
|
||||
let peerId = peermgr.switch.peerInfo.peerId
|
||||
var spamProtectionConfig = defaultConfig()
|
||||
spamProtectionConfig.keystorePath = "rln_keystore_" & $peerId & ".json"
|
||||
spamProtectionConfig.keystorePassword = "mix-rln-password"
|
||||
if userMessageLimit.isSome():
|
||||
spamProtectionConfig.userMessageLimit = userMessageLimit.get()
|
||||
# rlnResourcesPath left empty to use bundled resources (via "tree_height_/" placeholder)
|
||||
|
||||
let spamProtection = newMixRlnSpamProtection(spamProtectionConfig).valueOr:
|
||||
return err("failed to create spam protection: " & error)
|
||||
|
||||
var m = WakuMix(
|
||||
peerManager: peermgr,
|
||||
clusterId: clusterId,
|
||||
pubKey: mixPubKey,
|
||||
mixRlnSpamProtection: spamProtection,
|
||||
publishMessage: publishMessage,
|
||||
)
|
||||
procCall MixProtocol(m).init(
|
||||
localMixNodeInfo,
|
||||
peermgr.switch,
|
||||
spamProtection = Opt.some(SpamProtection(spamProtection)),
|
||||
delayStrategy = Opt.some(
|
||||
DelayStrategy(
|
||||
ExponentialDelayStrategy.new(meanDelay = 50'u16, rng = crypto.newRng())
|
||||
ExponentialDelayStrategy.new(meanDelay = 100, rng = crypto.newRng())
|
||||
)
|
||||
),
|
||||
)
|
||||
@ -103,9 +138,218 @@ proc new*(
|
||||
|
||||
if m.nodePool.len < minMixPoolSize:
|
||||
warn "publishing with mix won't work until atleast 3 mix nodes in node pool"
|
||||
|
||||
return ok(m)
|
||||
|
||||
proc poolSize*(mix: WakuMix): int =
|
||||
mix.nodePool.len
|
||||
|
||||
proc setupSpamProtectionCallbacks(mix: WakuMix) =
|
||||
## Set up the publish callback for spam protection coordination.
|
||||
## This enables the plugin to broadcast membership updates and proof metadata
|
||||
## via Waku relay.
|
||||
if mix.publishMessage.isNil():
|
||||
warn "PublishMessage callback not available, spam protection coordination disabled"
|
||||
return
|
||||
|
||||
let publishCallback: PublishCallback = proc(
|
||||
contentTopic: string, data: seq[byte]
|
||||
) {.async.} =
|
||||
# Create a WakuMessage for the coordination data
|
||||
let msg = WakuMessage(
|
||||
payload: data,
|
||||
contentTopic: contentTopic,
|
||||
ephemeral: true, # Coordination messages don't need to be stored
|
||||
timestamp: getNowInNanosecondTime(),
|
||||
)
|
||||
|
||||
# Delegate to node's publish API which handles topic derivation and relay publishing
|
||||
let res = await mix.publishMessage(msg)
|
||||
if res.isErr():
|
||||
warn "Failed to publish spam protection coordination message",
|
||||
contentTopic = contentTopic, error = res.error
|
||||
return
|
||||
|
||||
trace "Published spam protection coordination message", contentTopic = contentTopic
|
||||
|
||||
mix.mixRlnSpamProtection.setPublishCallback(publishCallback)
|
||||
trace "Spam protection publish callback configured"
|
||||
|
||||
proc handleMessage*(
|
||||
mix: WakuMix, pubsubTopic: PubsubTopic, message: WakuMessage
|
||||
) {.async, gcsafe.} =
|
||||
## Handle incoming messages for spam protection coordination.
|
||||
## This should be called from the relay handler for coordination content topics.
|
||||
if mix.mixRlnSpamProtection.isNil():
|
||||
return
|
||||
|
||||
let contentTopic = message.contentTopic
|
||||
|
||||
if contentTopic == mix.mixRlnSpamProtection.getMembershipContentTopic():
|
||||
# Handle membership update
|
||||
let res = await mix.mixRlnSpamProtection.handleMembershipUpdate(message.payload)
|
||||
if res.isErr:
|
||||
warn "Failed to handle membership update", error = res.error
|
||||
else:
|
||||
trace "Handled membership update"
|
||||
|
||||
# Persist tree after membership changes (temporary solution)
|
||||
# TODO: Replace with proper persistence strategy (e.g., periodic snapshots)
|
||||
let saveRes = mix.mixRlnSpamProtection.saveTree()
|
||||
if saveRes.isErr:
|
||||
debug "Failed to save tree after membership update", error = saveRes.error
|
||||
else:
|
||||
trace "Saved tree after membership update"
|
||||
elif contentTopic == mix.mixRlnSpamProtection.getProofMetadataContentTopic():
|
||||
# Handle proof metadata for network-wide spam detection
|
||||
let res = mix.mixRlnSpamProtection.handleProofMetadata(message.payload)
|
||||
if res.isErr:
|
||||
warn "Failed to handle proof metadata", error = res.error
|
||||
else:
|
||||
trace "Handled proof metadata"
|
||||
|
||||
proc getSpamProtectionContentTopics*(mix: WakuMix): seq[string] =
|
||||
## Get the content topics used by spam protection for coordination.
|
||||
## Use these to set up relay subscriptions.
|
||||
if mix.mixRlnSpamProtection.isNil():
|
||||
return @[]
|
||||
return mix.mixRlnSpamProtection.getContentTopics()
|
||||
|
||||
proc saveSpamProtectionTree*(mix: WakuMix): Result[void, string] =
|
||||
## Save the spam protection membership tree to disk.
|
||||
## This allows preserving the tree state across restarts.
|
||||
if mix.mixRlnSpamProtection.isNil():
|
||||
return err("Spam protection not initialized")
|
||||
|
||||
mix.mixRlnSpamProtection.saveTree().mapErr(
|
||||
proc(e: string): string =
|
||||
e
|
||||
)
|
||||
|
||||
proc loadSpamProtectionTree*(mix: WakuMix): Result[void, string] =
|
||||
## Load the spam protection membership tree from disk.
|
||||
## Call this before init() to restore tree state from previous runs.
|
||||
## TODO: This is a temporary solution. Ideally nodes should sync tree state
|
||||
## via a store query for historical membership messages or via dedicated
|
||||
## tree sync protocol.
|
||||
if mix.mixRlnSpamProtection.isNil():
|
||||
return err("Spam protection not initialized")
|
||||
|
||||
mix.mixRlnSpamProtection.loadTree().mapErr(
|
||||
proc(e: string): string =
|
||||
e
|
||||
)
|
||||
|
||||
method start*(mix: WakuMix) {.async.} =
|
||||
## Local-only mix protocol initialization. Does NOT touch the network.
|
||||
## The network-dependent self-registration broadcast is handled separately
|
||||
## by registerDoSProtectionWithNetwork so that this proc can run before
|
||||
## peers are connected without blocking on relay startup.
|
||||
info "starting waku mix protocol"
|
||||
|
||||
if not mix.mixRlnSpamProtection.isNil():
|
||||
# Initialize spam protection (MixProtocol.init() does NOT call init() on the plugin)
|
||||
let initRes = await mix.mixRlnSpamProtection.init()
|
||||
if initRes.isErr:
|
||||
error "Failed to initialize spam protection", error = initRes.error
|
||||
return
|
||||
|
||||
# Load existing tree to sync with other members.
|
||||
# Should be done after init() (which loads credentials) but before
|
||||
# registerSelf() (which adds us to the tree).
|
||||
let loadRes = mix.mixRlnSpamProtection.loadTree()
|
||||
if loadRes.isErr:
|
||||
debug "No existing tree found or failed to load, starting fresh",
|
||||
error = loadRes.error
|
||||
else:
|
||||
debug "Loaded existing spam protection membership tree from disk"
|
||||
|
||||
# Restore our credentials to the tree (after tree load, whether it succeeded or not).
|
||||
# Ensures our member is in the tree if we have an index from keystore.
|
||||
let restoreRes = mix.mixRlnSpamProtection.restoreCredentialsToTree()
|
||||
if restoreRes.isErr:
|
||||
error "Failed to restore credentials to tree", error = restoreRes.error
|
||||
|
||||
# Set up publish callback. Must be before the network-side registration so
|
||||
# the plugin's groupManager.register can broadcast the membership update.
|
||||
mix.setupSpamProtectionCallbacks()
|
||||
|
||||
let startRes = await mix.mixRlnSpamProtection.start()
|
||||
if startRes.isErr:
|
||||
error "Failed to start spam protection", error = startRes.error
|
||||
|
||||
info "waku mix protocol started"
|
||||
|
||||
proc dosRegistrationRetryLoop(mix: WakuMix) {.async.} =
|
||||
## Indefinitely retry the DoS-protection self-registration broadcast until
|
||||
## it succeeds (or this task is cancelled by WakuMix.stop()). For nodes that
|
||||
## already have a membership index in their keystore, registerSelf early-
|
||||
## returns and the loop exits on the first attempt. For fresh nodes, the
|
||||
## broadcast needs at least one relay peer subscribed to the membership
|
||||
## topic to land — this loop survives transient "no peers yet" failures.
|
||||
##
|
||||
## TODO: Remove once RLN membership moves on-chain. With on-chain membership
|
||||
## peers discover each other via the contract / a watcher rather than via a
|
||||
## pubsub broadcast, so the retry loop (and the whole publishCallback path
|
||||
## from registerSelf) becomes unnecessary.
|
||||
##
|
||||
## Retry pacing uses exponential backoff (5s, 10s, 20s, ..., capped at 5min)
|
||||
## so persistent misconfiguration — e.g., relay never available — degrades
|
||||
## to one log line every 5 minutes after the initial ramp instead of every
|
||||
## 5 seconds forever.
|
||||
const InitialRetryDelay = chronos.seconds(5)
|
||||
const MaxRetryDelay = chronos.minutes(5)
|
||||
var delay = InitialRetryDelay
|
||||
while true:
|
||||
try:
|
||||
let registerRes = await mix.mixRlnSpamProtection.registerSelf()
|
||||
if registerRes.isOk():
|
||||
debug "DoS-protection self-registration succeeded", index = registerRes.get()
|
||||
# Persist tree only after a successful register — for fresh nodes this
|
||||
# captures the new index; for keystore nodes it's a harmless no-op.
|
||||
let saveRes = mix.mixRlnSpamProtection.saveTree()
|
||||
if saveRes.isErr:
|
||||
warn "Failed to save spam protection tree", error = saveRes.error
|
||||
else:
|
||||
trace "Saved spam protection tree to disk"
|
||||
return # success — exit the loop
|
||||
warn "DoS-protection self-registration failed, retrying",
|
||||
error = registerRes.error, nextDelay = delay
|
||||
except CancelledError as e:
|
||||
debug "DoS-protection registration loop cancelled"
|
||||
raise e
|
||||
except CatchableError as e:
|
||||
warn "DoS-protection registration raised, retrying",
|
||||
error = e.msg, nextDelay = delay
|
||||
await sleepAsync(delay)
|
||||
delay = min(delay * 2, MaxRetryDelay)
|
||||
|
||||
proc registerDoSProtectionWithNetwork*(mix: WakuMix) =
|
||||
## Kick off an indefinite background task that broadcasts this node's
|
||||
## DoS-protection (RLN) membership registration to other mix nodes via
|
||||
## relay. Returns immediately so callers don't block on a possibly-slow
|
||||
## broadcast. The task is cancelled when WakuMix.stop() is called.
|
||||
if mix.mixRlnSpamProtection.isNil():
|
||||
return
|
||||
# Guard against kicking off the retry loop when the plugin isn't actually
|
||||
# usable (e.g., mix.start()'s init/start steps failed). Without this check
|
||||
# the loop would spin forever logging "Plugin not initialized" warnings.
|
||||
if not mix.mixRlnSpamProtection.isReady():
|
||||
warn "Skipping DoS-protection registration: plugin not ready"
|
||||
return
|
||||
# Re-call safety: don't spawn a second loop if one is still in flight.
|
||||
if not mix.dosRegistrationTask.isNil and not mix.dosRegistrationTask.finished:
|
||||
debug "DoS-protection registration already in progress, skipping"
|
||||
return
|
||||
mix.dosRegistrationTask = mix.dosRegistrationRetryLoop()
|
||||
|
||||
method stop*(mix: WakuMix) {.async.} =
|
||||
# Cancel the in-flight DoS-protection registration retry loop, if any
|
||||
if not mix.dosRegistrationTask.isNil and not mix.dosRegistrationTask.finished:
|
||||
await mix.dosRegistrationTask.cancelAndWait()
|
||||
# Stop spam protection
|
||||
if not mix.mixRlnSpamProtection.isNil():
|
||||
await mix.mixRlnSpamProtection.stop()
|
||||
debug "Spam protection stopped"
|
||||
|
||||
# Mix Protocol
|
||||
|
||||
29
nimble.lock
29
nimble.lock
@ -699,9 +699,9 @@
|
||||
}
|
||||
},
|
||||
"libp2p_mix": {
|
||||
"version": "0.1.0",
|
||||
"vcsRevision": "380513117d556bf8f70066f5e72a7fd74fe36ba6",
|
||||
"url": "https://github.com/logos-co/nim-libp2p-mix",
|
||||
"version": "#50c4ab4fa788a33eb12a0a2cecaa708873352b58",
|
||||
"vcsRevision": "50c4ab4fa788a33eb12a0a2cecaa708873352b58",
|
||||
"url": "https://github.com/logos-co/nim-libp2p-mix.git",
|
||||
"downloadMethod": "git",
|
||||
"dependencies": [
|
||||
"nim",
|
||||
@ -715,7 +715,28 @@
|
||||
"unittest2"
|
||||
],
|
||||
"checksums": {
|
||||
"sha1": "ccfb0f0160ac15ac970471964c730d57edacad91"
|
||||
"sha1": "3994284d7d7cb413f2830a91e20c1c3bb6e91158"
|
||||
}
|
||||
},
|
||||
"mix_rln_spam_protection": {
|
||||
"version": "#61ee3e5aacb6b224b70e164ef7d0a5714fe66b26",
|
||||
"vcsRevision": "61ee3e5aacb6b224b70e164ef7d0a5714fe66b26",
|
||||
"url": "https://github.com/logos-co/mix-rln-spam-protection-plugin.git",
|
||||
"downloadMethod": "git",
|
||||
"dependencies": [
|
||||
"nim",
|
||||
"results",
|
||||
"stew",
|
||||
"chronicles",
|
||||
"chronos",
|
||||
"nimcrypto",
|
||||
"secp256k1",
|
||||
"json_serialization",
|
||||
"libp2p",
|
||||
"libp2p_mix"
|
||||
],
|
||||
"checksums": {
|
||||
"sha1": "03d1e7d7663a3136b1799d3dc507ee4689511235"
|
||||
}
|
||||
}
|
||||
},
|
||||
|
||||
13
nix/deps.nix
13
nix/deps.nix
@ -313,9 +313,16 @@
|
||||
};
|
||||
|
||||
libp2p_mix = pkgs.fetchgit {
|
||||
url = "https://github.com/logos-co/nim-libp2p-mix";
|
||||
rev = "380513117d556bf8f70066f5e72a7fd74fe36ba6";
|
||||
sha256 = "05zjf98nl2hxx62m9blk4yip2f31y44r5x4n98lmm5hghb7wbcpk";
|
||||
url = "https://github.com/logos-co/nim-libp2p-mix.git";
|
||||
rev = "50c4ab4fa788a33eb12a0a2cecaa708873352b58";
|
||||
sha256 = "16prk6cqhalzsvh9kaif5cdn1yadssx3h4572j58fsgm20kdrala";
|
||||
fetchSubmodules = true;
|
||||
};
|
||||
|
||||
mix_rln_spam_protection = pkgs.fetchgit {
|
||||
url = "https://github.com/logos-co/mix-rln-spam-protection-plugin.git";
|
||||
rev = "61ee3e5aacb6b224b70e164ef7d0a5714fe66b26";
|
||||
sha256 = "0j68v3a8vwrrdpcfmabzdlx867nh4lf9flxvfrzq3xs03m2si57h";
|
||||
fetchSubmodules = true;
|
||||
};
|
||||
|
||||
|
||||
@ -3,66 +3,128 @@
|
||||
## Aim
|
||||
|
||||
Simulate a local mixnet along with a chat app to publish using mix.
|
||||
This is helpful to test any changes while development.
|
||||
It includes scripts that run a `4 node` mixnet along with a lightpush service node(without mix) that can be used to test quickly.
|
||||
This is helpful to test any changes during development.
|
||||
|
||||
## Simulation Details
|
||||
|
||||
Note that before running the simulation both `wakunode2` and `chat2mix` have to be built.
|
||||
The simulation includes:
|
||||
|
||||
1. A 5-node mixnet where `run_mix_node.sh` is the bootstrap node for the other 4 nodes
|
||||
2. Two chat app instances that publish messages using lightpush protocol over the mixnet
|
||||
|
||||
### Available Scripts
|
||||
|
||||
| Script | Description |
|
||||
| ------------------ | ------------------------------------------ |
|
||||
| `run_mix_node.sh` | Bootstrap mix node (must be started first) |
|
||||
| `run_mix_node1.sh` | Mix node 1 |
|
||||
| `run_mix_node2.sh` | Mix node 2 |
|
||||
| `run_mix_node3.sh` | Mix node 3 |
|
||||
| `run_mix_node4.sh` | Mix node 4 |
|
||||
| `run_chat_mix.sh` | Chat app instance 1 |
|
||||
| `run_chat_mix1.sh` | Chat app instance 2 |
|
||||
| `build_setup.sh` | Build and generate RLN credentials |
|
||||
|
||||
## Prerequisites
|
||||
|
||||
Before running the simulation, build `wakunode2` and `chat2mix`:
|
||||
|
||||
```bash
|
||||
cd <repo-root-dir>
|
||||
make wakunode2
|
||||
make chat2mix
|
||||
source env.sh
|
||||
make wakunode2 chat2mix
|
||||
```
|
||||
|
||||
Simulation includes scripts for:
|
||||
## RLN Spam Protection Setup
|
||||
|
||||
1. a 4 waku-node mixnet where `node1` is bootstrap node for the other 3 nodes.
|
||||
2. scripts to run chat app that publishes using lightpush protocol over the mixnet
|
||||
Generate RLN credentials and the shared Merkle tree for all nodes:
|
||||
|
||||
```bash
|
||||
cd simulations/mixnet
|
||||
./build_setup.sh
|
||||
```
|
||||
|
||||
This script will:
|
||||
|
||||
1. Build and run the `setup_credentials` tool
|
||||
2. Generate RLN credentials for all nodes (5 mix nodes + 2 chat clients)
|
||||
3. Create `rln_tree.db` - the shared Merkle tree with all members
|
||||
4. Create keystore files (`rln_keystore_{peerId}.json`) for each node
|
||||
|
||||
**Important:** All scripts must be run from this directory (`simulations/mixnet/`) so they can access their credentials and tree file.
|
||||
|
||||
To regenerate credentials (e.g., after adding new nodes), run `./build_setup.sh` again - it will clean up old files first.
|
||||
|
||||
## Usage
|
||||
|
||||
Start the service node with below command, which acts as bootstrap node for all other mix nodes.
|
||||
### Step 1: Start the Mix Nodes
|
||||
|
||||
`./run_lp_service_node.sh`
|
||||
Start the bootstrap node first (in a separate terminal):
|
||||
|
||||
To run the nodes for mixnet run the 4 node scripts in different terminals as below.
|
||||
|
||||
`./run_mix_node1.sh`
|
||||
|
||||
Look for following 2 log lines to ensure node ran successfully and has also mounted mix protocol.
|
||||
|
||||
```log
|
||||
INF 2025-08-01 14:51:05.445+05:30 mounting mix protocol topics="waku node" tid=39996871 file=waku_node.nim:231 nodeId="(listenAddresses: @[\"/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o\"], enrUri: \"enr:-NC4QKYtas8STkenlqBTJ3a1TTLzJA2DsGGbFlnxem9aSM2IXm-CSVZULdk2467bAyFnepnt8KP_QlfDzdaMXd_zqtwBgmlkgnY0gmlwhH8AAAGHbWl4LWtleaCdCc5iT3bo9gYmXtucyit96bQXcqbXhL3a-S_6j7p9LIptdWx0aWFkZHJzgIJyc4UAAgEAAIlzZWNwMjU2azGhA6RFtVJVBh0SYOoP8xrgnXSlpiFARmQkF9d8Rn4fSeiog3RjcILqYYN1ZHCCIymFd2FrdTIt\")"
|
||||
|
||||
INF 2025-08-01 14:49:23.467+05:30 Node setup complete topics="wakunode main" tid=39994244 file=wakunode2.nim:104
|
||||
```bash
|
||||
./run_mix_node.sh
|
||||
```
|
||||
|
||||
Once all the 4 nodes are up without any issues, run the script to start the chat application.
|
||||
Look for the following log lines to ensure the node started successfully:
|
||||
|
||||
`./run_chat_app.sh`
|
||||
```log
|
||||
INF mounting mix protocol topics="waku node"
|
||||
INF Node setup complete topics="wakunode main"
|
||||
```
|
||||
|
||||
Enter a nickname to be used.
|
||||
Verify RLN spam protection initialized correctly by checking for these logs:
|
||||
|
||||
```log
|
||||
INF Initializing MixRlnSpamProtection
|
||||
INF MixRlnSpamProtection initialized, waiting for sync
|
||||
DBG Tree loaded from file
|
||||
INF MixRlnSpamProtection started
|
||||
```
|
||||
|
||||
Then start the remaining mix nodes in separate terminals:
|
||||
|
||||
```bash
|
||||
./run_mix_node1.sh
|
||||
./run_mix_node2.sh
|
||||
./run_mix_node3.sh
|
||||
./run_mix_node4.sh
|
||||
```
|
||||
|
||||
### Step 2: Start the Chat Applications
|
||||
|
||||
Once all 5 mix nodes are running, start the first chat app:
|
||||
|
||||
```bash
|
||||
./run_chat_mix.sh
|
||||
```
|
||||
|
||||
Enter a nickname when prompted:
|
||||
|
||||
```bash
|
||||
pubsub topic is: /waku/2/rs/2/0
|
||||
Choose a nickname >>
|
||||
```
|
||||
|
||||
Once you see below log, it means the app is ready for publishing messages over the mixnet.
|
||||
Once you see the following log, the app is ready to publish messages over the mixnet:
|
||||
|
||||
```bash
|
||||
Welcome, test!
|
||||
Listening on
|
||||
/ip4/192.168.68.64/tcp/60000/p2p/16Uiu2HAkxDGqix1ifY3wF1ZzojQWRAQEdKP75wn1LJMfoHhfHz57
|
||||
/ip4/<local-network-ip>/tcp/60000/p2p/16Uiu2HAkxDGqix1ifY3wF1ZzojQWRAQEdKP75wn1LJMfoHhfHz57
|
||||
ready to publish messages now
|
||||
```
|
||||
|
||||
Follow similar instructions to run second instance of chat app.
|
||||
Once both the apps run successfully, send a message and check if it is received by the other app.
|
||||
Start the second chat app in another terminal:
|
||||
|
||||
You can exit the chat apps by entering `/exit` as below
|
||||
```bash
|
||||
./run_chat_mix1.sh
|
||||
```
|
||||
|
||||
### Step 3: Test Messaging
|
||||
|
||||
Once both chat apps are running, send a message from one and verify it is received by the other.
|
||||
|
||||
To exit the chat apps, enter `/exit`:
|
||||
|
||||
```bash
|
||||
>> /exit
|
||||
|
||||
36
simulations/mixnet/build_setup.sh
Executable file
36
simulations/mixnet/build_setup.sh
Executable file
@ -0,0 +1,36 @@
|
||||
#!/bin/bash
|
||||
cd "$(dirname "$0")"
|
||||
MIXNET_DIR=$(pwd)
|
||||
cd ../..
|
||||
ROOT_DIR=$(pwd)
|
||||
source "$ROOT_DIR/env.sh"
|
||||
|
||||
# Prefer explicitly provided RLN library path, otherwise use the one built by `make librln`.
|
||||
LIBRLN_PATH=${LIBRLN_PATH:-"$ROOT_DIR/librln_v2.0.2.a"}
|
||||
|
||||
# Clean up old files first
|
||||
rm -f "$MIXNET_DIR/rln_tree.db" "$MIXNET_DIR"/rln_keystore_*.json
|
||||
|
||||
echo "Building and running credentials setup..."
|
||||
# Compile to temp location, then run from mixnet directory
|
||||
nim c -d:release --mm:refc \
|
||||
--passL:"$LIBRLN_PATH" --passL:-lm \
|
||||
-o:/tmp/setup_credentials_$$ \
|
||||
"$MIXNET_DIR/setup_credentials.nim" 2>&1 | tail -30
|
||||
|
||||
# Run from mixnet directory so files are created there
|
||||
cd "$MIXNET_DIR"
|
||||
/tmp/setup_credentials_$$
|
||||
|
||||
# Clean up temp binary
|
||||
rm -f /tmp/setup_credentials_$$
|
||||
|
||||
# Verify output
|
||||
if [ -f "rln_tree.db" ]; then
|
||||
echo ""
|
||||
echo "Tree file ready at: $(pwd)/rln_tree.db"
|
||||
ls -la rln_keystore_*.json 2>/dev/null | wc -l | xargs -I {} echo "Generated {} keystore files"
|
||||
else
|
||||
echo "Setup failed - rln_tree.db not found"
|
||||
exit 1
|
||||
fi
|
||||
@ -13,7 +13,7 @@ discv5-udp-port = 9002
|
||||
discv5-enr-auto-update = true
|
||||
discv5-bootstrap-node = ["enr:-LG4QBaAbcA921hmu3IrreLqGZ4y3VWCjBCgNN9mpX9vqkkbSrM3HJHZTXnb5iVXgc5pPtDhWLxkB6F3yY25hSwMezkEgmlkgnY0gmlwhH8AAAGKbXVsdGlhZGRyc4oACATAqEQ-BuphgnJzhQACAQAAiXNlY3AyNTZrMaEDpEW1UlUGHRJg6g_zGuCddKWmIUBGZCQX13xGfh9J6KiDdGNwguphg3VkcIIjKYV3YWt1Mg0"]
|
||||
kad-bootstrap-node = ["/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o"]
|
||||
rest = false
|
||||
rest = true
|
||||
rest-admin = false
|
||||
ports-shift = 3
|
||||
num-shards-in-network = 1
|
||||
|
||||
@ -13,7 +13,7 @@ discv5-udp-port = 9003
|
||||
discv5-enr-auto-update = true
|
||||
discv5-bootstrap-node = ["enr:-LG4QBaAbcA921hmu3IrreLqGZ4y3VWCjBCgNN9mpX9vqkkbSrM3HJHZTXnb5iVXgc5pPtDhWLxkB6F3yY25hSwMezkEgmlkgnY0gmlwhH8AAAGKbXVsdGlhZGRyc4oACATAqEQ-BuphgnJzhQACAQAAiXNlY3AyNTZrMaEDpEW1UlUGHRJg6g_zGuCddKWmIUBGZCQX13xGfh9J6KiDdGNwguphg3VkcIIjKYV3YWt1Mg0"]
|
||||
kad-bootstrap-node = ["/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o"]
|
||||
rest = false
|
||||
rest = true
|
||||
rest-admin = false
|
||||
ports-shift = 4
|
||||
num-shards-in-network = 1
|
||||
|
||||
@ -13,7 +13,7 @@ discv5-udp-port = 9004
|
||||
discv5-enr-auto-update = true
|
||||
discv5-bootstrap-node = ["enr:-LG4QBaAbcA921hmu3IrreLqGZ4y3VWCjBCgNN9mpX9vqkkbSrM3HJHZTXnb5iVXgc5pPtDhWLxkB6F3yY25hSwMezkEgmlkgnY0gmlwhH8AAAGKbXVsdGlhZGRyc4oACATAqEQ-BuphgnJzhQACAQAAiXNlY3AyNTZrMaEDpEW1UlUGHRJg6g_zGuCddKWmIUBGZCQX13xGfh9J6KiDdGNwguphg3VkcIIjKYV3YWt1Mg0"]
|
||||
kad-bootstrap-node = ["/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o"]
|
||||
rest = false
|
||||
rest = true
|
||||
rest-admin = false
|
||||
ports-shift = 5
|
||||
num-shards-in-network = 1
|
||||
|
||||
@ -1,2 +1,2 @@
|
||||
../../build/chat2mix --cluster-id=2 --num-shards-in-network=1 --shard=0 --servicenode="/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o" --log-level=TRACE --kad-bootstrap-node="/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o"
|
||||
../../build/chat2mix --cluster-id=2 --num-shards-in-network=1 --shard=0 --servicenode="/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o" --log-level=TRACE --nodekey="cb6fe589db0e5d5b48f7e82d33093e4d9d35456f4aaffc2322c473a173b2ac49" --kad-bootstrap-node="/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o" --fleet="none"
|
||||
#--mixnode="/ip4/127.0.0.1/tcp/60002/p2p/16Uiu2HAmLtKaFaSWDohToWhWUZFLtqzYZGPFuXwKrojFVF6az5UF:9231e86da6432502900a84f867004ce78632ab52cd8e30b1ec322cd795710c2a" --mixnode="/ip4/127.0.0.1/tcp/60003/p2p/16Uiu2HAmTEDHwAziWUSz6ZE23h5vxG2o4Nn7GazhMor4bVuMXTrA:275cd6889e1f29ca48e5b9edb800d1a94f49f13d393a0ecf1a07af753506de6c" --mixnode="/ip4/127.0.0.1/tcp/60004/p2p/16Uiu2HAmPwRKZajXtfb1Qsv45VVfRZgK3ENdfmnqzSrVm3BczF6f:e0ed594a8d506681be075e8e23723478388fb182477f7a469309a25e7076fc18" --mixnode="/ip4/127.0.0.1/tcp/60005/p2p/16Uiu2HAmRhxmCHBYdXt1RibXrjAUNJbduAhzaTHwFCZT4qWnqZAu:8fd7a1a7c19b403d231452a9b1ea40eb1cc76f455d918ef8980e7685f9eeeb1f"
|
||||
|
||||
@ -1,2 +1 @@
|
||||
../../build/chat2mix --cluster-id=2 --num-shards-in-network=1 --shard=0 --servicenode="/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o" --log-level=TRACE
|
||||
#--mixnode="/ip4/127.0.0.1/tcp/60002/p2p/16Uiu2HAmLtKaFaSWDohToWhWUZFLtqzYZGPFuXwKrojFVF6az5UF:9231e86da6432502900a84f867004ce78632ab52cd8e30b1ec322cd795710c2a" --mixnode="/ip4/127.0.0.1/tcp/60003/p2p/16Uiu2HAmTEDHwAziWUSz6ZE23h5vxG2o4Nn7GazhMor4bVuMXTrA:275cd6889e1f29ca48e5b9edb800d1a94f49f13d393a0ecf1a07af753506de6c" --mixnode="/ip4/127.0.0.1/tcp/60004/p2p/16Uiu2HAmPwRKZajXtfb1Qsv45VVfRZgK3ENdfmnqzSrVm3BczF6f:e0ed594a8d506681be075e8e23723478388fb182477f7a469309a25e7076fc18" --mixnode="/ip4/127.0.0.1/tcp/60005/p2p/16Uiu2HAmRhxmCHBYdXt1RibXrjAUNJbduAhzaTHwFCZT4qWnqZAu:8fd7a1a7c19b403d231452a9b1ea40eb1cc76f455d918ef8980e7685f9eeeb1f"
|
||||
../../build/chat2mix --cluster-id=2 --num-shards-in-network=1 --shard=0 --servicenode="/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o" --log-level=TRACE --nodekey="35eace7ccb246f20c487e05015ca77273d8ecaed0ed683de3d39bf4f69336feb" --mixnode="/ip4/127.0.0.1/tcp/60002/p2p/16Uiu2HAmLtKaFaSWDohToWhWUZFLtqzYZGPFuXwKrojFVF6az5UF:9231e86da6432502900a84f867004ce78632ab52cd8e30b1ec322cd795710c2a" --mixnode="/ip4/127.0.0.1/tcp/60003/p2p/16Uiu2HAmTEDHwAziWUSz6ZE23h5vxG2o4Nn7GazhMor4bVuMXTrA:275cd6889e1f29ca48e5b9edb800d1a94f49f13d393a0ecf1a07af753506de6c" --mixnode="/ip4/127.0.0.1/tcp/60004/p2p/16Uiu2HAmPwRKZajXtfb1Qsv45VVfRZgK3ENdfmnqzSrVm3BczF6f:e0ed594a8d506681be075e8e23723478388fb182477f7a469309a25e7076fc18" --mixnode="/ip4/127.0.0.1/tcp/60005/p2p/16Uiu2HAmRhxmCHBYdXt1RibXrjAUNJbduAhzaTHwFCZT4qWnqZAu:8fd7a1a7c19b403d231452a9b1ea40eb1cc76f455d918ef8980e7685f9eeeb1f" --mixnode="/ip4/127.0.0.1/tcp/60001/p2p/16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o:9d09ce624f76e8f606265edb9cca2b7de9b41772a6d784bddaf92ffa8fba7d2c" --fleet="none"
|
||||
|
||||
@ -1 +1,2 @@
|
||||
../../build/wakunode2 --config-file="config.toml" 2>&1 | tee mix_node.log
|
||||
|
||||
|
||||
139
simulations/mixnet/setup_credentials.nim
Normal file
139
simulations/mixnet/setup_credentials.nim
Normal file
@ -0,0 +1,139 @@
|
||||
{.push raises: [].}
|
||||
|
||||
## Setup script to generate RLN credentials and shared Merkle tree for mix nodes.
|
||||
##
|
||||
## This script:
|
||||
## 1. Generates credentials for each node (identified by peer ID)
|
||||
## 2. Registers all credentials in a shared Merkle tree
|
||||
## 3. Saves the tree to rln_tree.db
|
||||
## 4. Saves individual keystores named by peer ID
|
||||
##
|
||||
## Usage: nim c -r setup_credentials.nim
|
||||
|
||||
import std/[os, strformat, options], chronicles, chronos, results
|
||||
|
||||
import
|
||||
mix_rln_spam_protection/credentials,
|
||||
mix_rln_spam_protection/group_manager,
|
||||
mix_rln_spam_protection/rln_interface,
|
||||
mix_rln_spam_protection/types
|
||||
|
||||
const
|
||||
KeystorePassword = "mix-rln-password" # Must match protocol.nim
|
||||
DefaultUserMessageLimit = 100'u64 # Network-wide default rate limit
|
||||
SpammerUserMessageLimit = 3'u64 # Lower limit for spammer testing
|
||||
|
||||
# Peer IDs derived from nodekeys in config files
|
||||
# config.toml: nodekey = "f98e3fba96c32e8d1967d460f1b79457380e1a895f7971cecc8528abe733781a"
|
||||
# config1.toml: nodekey = "09e9d134331953357bd38bbfce8edb377f4b6308b4f3bfbe85c610497053d684"
|
||||
# config2.toml: nodekey = "ed54db994682e857d77cd6fb81be697382dc43aa5cd78e16b0ec8098549f860e"
|
||||
# config3.toml: nodekey = "42f96f29f2d6670938b0864aced65a332dcf5774103b4c44ec4d0ea4ef3c47d6"
|
||||
# config4.toml: nodekey = "3ce887b3c34b7a92dd2868af33941ed1dbec4893b054572cd5078da09dd923d4"
|
||||
# chat2mix.sh: nodekey = "cb6fe589db0e5d5b48f7e82d33093e4d9d35456f4aaffc2322c473a173b2ac49"
|
||||
# chat2mix1.sh: nodekey = "35eace7ccb246f20c487e05015ca77273d8ecaed0ed683de3d39bf4f69336feb"
|
||||
|
||||
# Node info: (peerId, userMessageLimit)
|
||||
NodeConfigs = [
|
||||
("16Uiu2HAmPiEs2ozjjJF2iN2Pe2FYeMC9w4caRHKYdLdAfjgbWM6o", DefaultUserMessageLimit),
|
||||
# config.toml (service node)
|
||||
("16Uiu2HAmLtKaFaSWDohToWhWUZFLtqzYZGPFuXwKrojFVF6az5UF", DefaultUserMessageLimit),
|
||||
# config1.toml (mix node 1)
|
||||
("16Uiu2HAmTEDHwAziWUSz6ZE23h5vxG2o4Nn7GazhMor4bVuMXTrA", DefaultUserMessageLimit),
|
||||
# config2.toml (mix node 2)
|
||||
("16Uiu2HAmPwRKZajXtfb1Qsv45VVfRZgK3ENdfmnqzSrVm3BczF6f", DefaultUserMessageLimit),
|
||||
# config3.toml (mix node 3)
|
||||
("16Uiu2HAmRhxmCHBYdXt1RibXrjAUNJbduAhzaTHwFCZT4qWnqZAu", DefaultUserMessageLimit),
|
||||
# config4.toml (mix node 4)
|
||||
("16Uiu2HAm1QxSjNvNbsT2xtLjRGAsBLVztsJiTHr9a3EK96717hpj", DefaultUserMessageLimit),
|
||||
# chat2mix client 1
|
||||
("16Uiu2HAmC9h26U1C83FJ5xpE32ghqya8CaZHX1Y7qpfHNnRABscN", DefaultUserMessageLimit),
|
||||
# chat2mix client 2
|
||||
]
|
||||
|
||||
proc setupCredentialsAndTree() {.async.} =
|
||||
## Generate credentials for all nodes and create a shared tree
|
||||
|
||||
echo "=== RLN Credentials Setup ==="
|
||||
echo "Generating credentials for ", NodeConfigs.len, " nodes...\n"
|
||||
|
||||
# Generate credentials for all nodes
|
||||
var allCredentials:
|
||||
seq[tuple[peerId: string, cred: IdentityCredential, rateLimit: uint64]]
|
||||
for (peerId, rateLimit) in NodeConfigs:
|
||||
let cred = generateCredentials().valueOr:
|
||||
echo "Failed to generate credentials for ", peerId, ": ", error
|
||||
quit(1)
|
||||
|
||||
allCredentials.add((peerId: peerId, cred: cred, rateLimit: rateLimit))
|
||||
echo "Generated credentials for ", peerId
|
||||
echo " idCommitment: ", cred.idCommitment.toHex()[0 .. 15], "..."
|
||||
echo " userMessageLimit: ", rateLimit
|
||||
|
||||
echo ""
|
||||
|
||||
# Create a group manager directly to build the tree
|
||||
let rlnInstance = newRLNInstance().valueOr:
|
||||
echo "Failed to create RLN instance: ", error
|
||||
quit(1)
|
||||
|
||||
let groupManager = newOffchainGroupManager(rlnInstance, "/mix/rln/membership/v1")
|
||||
|
||||
# Initialize the group manager
|
||||
let initRes = await groupManager.init()
|
||||
if initRes.isErr:
|
||||
echo "Failed to initialize group manager: ", initRes.error
|
||||
quit(1)
|
||||
|
||||
# Register all credentials in the tree with their specific rate limits
|
||||
echo "Registering all credentials in the Merkle tree..."
|
||||
for i, entry in allCredentials:
|
||||
let index = (
|
||||
await groupManager.registerWithLimit(entry.cred.idCommitment, entry.rateLimit)
|
||||
).valueOr:
|
||||
echo "Failed to register credential for ", entry.peerId, ": ", error
|
||||
quit(1)
|
||||
echo " Registered ",
|
||||
entry.peerId, " at index ", index, " (limit: ", entry.rateLimit, ")"
|
||||
|
||||
echo ""
|
||||
|
||||
# Save the tree to disk
|
||||
echo "Saving tree to rln_tree.db..."
|
||||
let saveRes = groupManager.saveTreeToFile("rln_tree.db")
|
||||
if saveRes.isErr:
|
||||
echo "Failed to save tree: ", saveRes.error
|
||||
quit(1)
|
||||
echo "Tree saved successfully!"
|
||||
|
||||
echo ""
|
||||
|
||||
# Save each credential to a keystore file named by peer ID
|
||||
echo "Saving keystores..."
|
||||
for i, entry in allCredentials:
|
||||
let keystorePath = &"rln_keystore_{entry.peerId}.json"
|
||||
|
||||
# Save with membership index and rate limit
|
||||
let saveResult = saveKeystore(
|
||||
entry.cred,
|
||||
KeystorePassword,
|
||||
keystorePath,
|
||||
some(MembershipIndex(i)),
|
||||
some(entry.rateLimit),
|
||||
)
|
||||
if saveResult.isErr:
|
||||
echo "Failed to save keystore for ", entry.peerId, ": ", saveResult.error
|
||||
quit(1)
|
||||
echo " Saved: ", keystorePath, " (limit: ", entry.rateLimit, ")"
|
||||
|
||||
echo ""
|
||||
echo "=== Setup Complete ==="
|
||||
echo " Tree file: rln_tree.db (", NodeConfigs.len, " members)"
|
||||
echo " Keystores: rln_keystore_{peerId}.json"
|
||||
echo " Password: ", KeystorePassword
|
||||
echo " Default rate limit: ", DefaultUserMessageLimit
|
||||
echo " Spammer rate limit: ", SpammerUserMessageLimit
|
||||
echo ""
|
||||
echo "Note: All nodes must use the same rln_tree.db file."
|
||||
|
||||
when isMainModule:
|
||||
waitFor setupCredentialsAndTree()
|
||||
@ -22,7 +22,7 @@ proc setupTestNode*(
|
||||
addAllCapabilities = false,
|
||||
bindUdpPort = address.udpPort, # Assume same as external
|
||||
bindTcpPort = address.tcpPort, # Assume same as external
|
||||
rng = rng,
|
||||
rng = rng(),
|
||||
)
|
||||
nextPort.inc
|
||||
for capability in capabilities:
|
||||
@ -40,7 +40,10 @@ proc getRng(): crypto.Rng =
|
||||
# purpose of the tests, it's ok as long as we only use a single thread
|
||||
{.gcsafe.}:
|
||||
if rngVar.rng.isNil:
|
||||
rngVar.rng = crypto.newRng()
|
||||
# libp2p v2.0.0: crypto.newRng() returns the new `Rng` wrapper type;
|
||||
# construct an HmacDrbgContext directly so the field type stays as
|
||||
# `ref HmacDrbgContext` (what bearssl-style consumers expect).
|
||||
rngVar.rng = HmacDrbgContext.new()
|
||||
rngVar.rng
|
||||
|
||||
template rng*(): crypto.Rng =
|
||||
|
||||
@ -1346,7 +1346,8 @@ procSuite "Peer Manager":
|
||||
|
||||
# Create peer manager
|
||||
let pm = PeerManager.new(
|
||||
switch = SwitchBuilder.new().withRng(rng()).withMplex().withNoise().build(),
|
||||
switch =
|
||||
SwitchBuilder.new().withRng(crypto.newRng()).withMplex().withNoise().build(),
|
||||
storage = nil,
|
||||
)
|
||||
|
||||
|
||||
@ -4,6 +4,7 @@ import
|
||||
testutils/unittests,
|
||||
chronos,
|
||||
libp2p/builders,
|
||||
libp2p/crypto/crypto,
|
||||
libp2p/protocols/connectivity/autonat/client,
|
||||
libp2p/protocols/connectivity/relay/relay,
|
||||
libp2p/protocols/connectivity/relay/client,
|
||||
@ -13,7 +14,7 @@ import logos_delivery/waku/node/waku_switch, ./testlib/common, ./testlib/wakucor
|
||||
proc newCircuitRelayClientSwitch(relayClient: RelayClient): Switch =
|
||||
SwitchBuilder
|
||||
.new()
|
||||
.withRng(rng())
|
||||
.withRng(crypto.newRng())
|
||||
.withAddresses(@[MultiAddress.init("/ip4/0.0.0.0/tcp/0").tryGet()])
|
||||
.withTcpTransport()
|
||||
.withMplex()
|
||||
@ -26,7 +27,7 @@ suite "Waku Switch":
|
||||
## Given
|
||||
let
|
||||
sourceSwitch = newTestSwitch()
|
||||
wakuSwitch = newWakuSwitch(rng = rng(), circuitRelay = Relay.new())
|
||||
wakuSwitch = newWakuSwitch(rng = crypto.newRng(), circuitRelay = Relay.new())
|
||||
await sourceSwitch.start()
|
||||
await wakuSwitch.start()
|
||||
|
||||
@ -46,7 +47,7 @@ suite "Waku Switch":
|
||||
asyncTest "Waku Switch acts as circuit relayer":
|
||||
## Setup
|
||||
let
|
||||
wakuSwitch = newWakuSwitch(rng = rng(), circuitRelay = Relay.new())
|
||||
wakuSwitch = newWakuSwitch(rng = crypto.newRng(), circuitRelay = Relay.new())
|
||||
sourceClient = RelayClient.new()
|
||||
destClient = RelayClient.new()
|
||||
sourceSwitch = newCircuitRelayClientSwitch(sourceClient)
|
||||
|
||||
@ -39,8 +39,8 @@ suite "Waku Filter - End to End":
|
||||
pubsubTopic = DefaultPubsubTopic
|
||||
contentTopic = DefaultContentTopic
|
||||
contentTopicSeq = @[contentTopic]
|
||||
serverSwitch = newStandardSwitch()
|
||||
clientSwitch = newStandardSwitch()
|
||||
serverSwitch = newTestSwitch()
|
||||
clientSwitch = newTestSwitch()
|
||||
wakuFilter = await newTestWakuFilter(serverSwitch)
|
||||
wakuFilterClient = await newTestWakuFilterClient(clientSwitch)
|
||||
|
||||
@ -106,7 +106,7 @@ suite "Waku Filter - End to End":
|
||||
suite "Subscribe":
|
||||
asyncTest "Server remote peer info doesn't match an online server":
|
||||
# Given an offline service node
|
||||
let offlineServerSwitch = newStandardSwitch()
|
||||
let offlineServerSwitch = newTestSwitch()
|
||||
let offlineServerRemotePeerInfo =
|
||||
offlineServerSwitch.peerInfo.toRemotePeerInfo()
|
||||
|
||||
@ -721,7 +721,7 @@ suite "Waku Filter - End to End":
|
||||
# Given a WakuFilterClient list of size MaxFilterPeers
|
||||
var clients: seq[(WakuFilterClient, Switch)] = @[]
|
||||
for i in 0 ..< MaxFilterPeers:
|
||||
let standardSwitch = newStandardSwitch()
|
||||
let standardSwitch = newTestSwitch()
|
||||
let wakuFilterClient = await newTestWakuFilterClient(standardSwitch)
|
||||
clients.add((wakuFilterClient, standardSwitch))
|
||||
|
||||
@ -738,7 +738,7 @@ suite "Waku Filter - End to End":
|
||||
wakuFilter.subscriptions.subscribedPeerCount() == MaxFilterPeers
|
||||
|
||||
# When initialising a new WakuFilterClient and subscribing it to the same service
|
||||
let standardSwitch = newStandardSwitch()
|
||||
let standardSwitch = newTestSwitch()
|
||||
let wakuFilterClient = await newTestWakuFilterClient(standardSwitch)
|
||||
await standardSwitch.start()
|
||||
let subscribeResponse = await wakuFilterClient.subscribe(
|
||||
@ -752,7 +752,7 @@ suite "Waku Filter - End to End":
|
||||
|
||||
asyncTest "Multiple Subscriptions":
|
||||
# Given a second service node
|
||||
let serverSwitch2 = newStandardSwitch()
|
||||
let serverSwitch2 = newTestSwitch()
|
||||
let wakuFilter2 = await newTestWakuFilter(serverSwitch2)
|
||||
await allFutures(serverSwitch2.start())
|
||||
let serverRemotePeerInfo2 = serverSwitch2.peerInfo.toRemotePeerInfo()
|
||||
|
||||
@ -24,7 +24,7 @@ type AFilterClient = ref object of RootObj
|
||||
|
||||
proc init(T: type[AFilterClient]): T =
|
||||
var r = T(
|
||||
clientSwitch: newStandardSwitch(),
|
||||
clientSwitch: newTestSwitch(),
|
||||
msgSeq: @[],
|
||||
pushHandlerFuture: newPushHandlerFuture(),
|
||||
)
|
||||
@ -93,7 +93,7 @@ suite "Waku Filter - DOS protection":
|
||||
pubsubTopic = DefaultPubsubTopic
|
||||
contentTopic = DefaultContentTopic
|
||||
contentTopicSeq = @[contentTopic]
|
||||
serverSwitch = newStandardSwitch()
|
||||
serverSwitch = newTestSwitch()
|
||||
wakuFilter = await newTestWakuFilter(
|
||||
serverSwitch, rateLimitSetting = some((3, 1000.milliseconds))
|
||||
)
|
||||
|
||||
@ -107,7 +107,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
peerSwitch.mount(node.wakuStore)
|
||||
@ -184,7 +184,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
peerSwitch.mount(node.wakuStore)
|
||||
@ -253,7 +253,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
peerSwitch.mount(node.wakuStore)
|
||||
@ -348,7 +348,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
peerSwitch.mount(node.wakuStore)
|
||||
@ -421,7 +421,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
peerSwitch.mount(node.wakuStore)
|
||||
@ -510,7 +510,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
peerSwitch.mount(node.wakuStore)
|
||||
@ -560,7 +560,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
let client = newRestHttpClient(initTAddress(restAddress, restPort))
|
||||
@ -742,7 +742,7 @@ procSuite "Waku Rest API - Store v3":
|
||||
node.mountStoreClient()
|
||||
|
||||
let key = generateEcdsaKey()
|
||||
var peerSwitch = newStandardSwitch(Opt.some(key))
|
||||
var peerSwitch = newTestSwitch(some(key))
|
||||
await peerSwitch.start()
|
||||
|
||||
peerSwitch.mount(node.wakuStore)
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user