diff --git a/.gitignore b/.gitignore index bfd8f269f..76c777297 100644 --- a/.gitignore +++ b/.gitignore @@ -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 diff --git a/apps/chat2mix/chat2mix.nim b/apps/chat2mix/chat2mix.nim index 155f1e855..51baf7806 100644 --- a/apps/chat2mix/chat2mix.nim +++ b/apps/chat2mix/chat2mix.nim @@ -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() diff --git a/apps/chat2mix/config_chat2mix.nim b/apps/chat2mix/config_chat2mix.nim index f77a729f4..de66e1b7d 100644 --- a/apps/chat2mix/config_chat2mix.nim +++ b/apps/chat2mix/config_chat2mix.nim @@ -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: diff --git a/config.nims b/config.nims index aaba3854e..25a49ddb9 100644 --- a/config.nims +++ b/config.nims @@ -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 diff --git a/logos_delivery.nimble b/logos_delivery.nimble index efe118595..101de1904 100644 --- a/logos_delivery.nimble +++ b/logos_delivery.nimble @@ -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 diff --git a/logos_delivery/waku/factory/waku.nim b/logos_delivery/waku/factory/waku.nim index e15a92fd0..061c19d4a 100644 --- a/logos_delivery/waku/factory/waku.nim +++ b/logos_delivery/waku/factory/waku.nim @@ -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 diff --git a/logos_delivery/waku/node/peer_manager/peer_manager.nim b/logos_delivery/waku/node/peer_manager/peer_manager.nim index 60bd2acb1..72f81e562 100644 --- a/logos_delivery/waku/node/peer_manager/peer_manager.nim +++ b/logos_delivery/waku/node/peer_manager/peer_manager.nim @@ -28,6 +28,7 @@ import node/health_monitor/online_monitor, node/waku_switch, ], + ../waku_switch, ./peer_store/peer_storage, ./waku_peer_store diff --git a/logos_delivery/waku/node/subscription_manager.nim b/logos_delivery/waku/node/subscription_manager.nim index 15b582ea6..58acd06aa 100644 --- a/logos_delivery/waku/node/subscription_manager.nim +++ b/logos_delivery/waku/node/subscription_manager.nim @@ -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(): diff --git a/logos_delivery/waku/node/waku_mix_coordination.nim b/logos_delivery/waku/node/waku_mix_coordination.nim new file mode 100644 index 000000000..3eb53750e --- /dev/null +++ b/logos_delivery/waku/node/waku_mix_coordination.nim @@ -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) diff --git a/logos_delivery/waku/node/waku_node.nim b/logos_delivery/waku/node/waku_node.nim index 2ad7dc601..74052cd62 100644 --- a/logos_delivery/waku/node/waku_node.nim +++ b/logos_delivery/waku/node/waku_node.nim @@ -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(): diff --git a/logos_delivery/waku/node/waku_node/relay.nim b/logos_delivery/waku/node/waku_node/relay.nim index 57904dc94..69619e860 100644 --- a/logos_delivery/waku/node/waku_node/relay.nim +++ b/logos_delivery/waku/node/waku_node/relay.nim @@ -29,6 +29,7 @@ import waku_archive, waku_store_sync, waku_rln_relay, + waku_mix, node/waku_node, node/subscription_manager, node/peer_manager, diff --git a/logos_delivery/waku/waku_mix/protocol.nim b/logos_delivery/waku/waku_mix/protocol.nim index 613b6e3c7..58b8b66d8 100644 --- a/logos_delivery/waku/waku_mix/protocol.nim +++ b/logos_delivery/waku/waku_mix/protocol.nim @@ -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 diff --git a/nimble.lock b/nimble.lock index 18ebde258..2d8aab5eb 100644 --- a/nimble.lock +++ b/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" } } }, diff --git a/nix/deps.nix b/nix/deps.nix index 00dea27c2..a447be7d6 100644 --- a/nix/deps.nix +++ b/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; }; diff --git a/simulations/mixnet/README.md b/simulations/mixnet/README.md index fcc67b6e1..99b0ba50b 100644 --- a/simulations/mixnet/README.md +++ b/simulations/mixnet/README.md @@ -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 -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//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 diff --git a/simulations/mixnet/build_setup.sh b/simulations/mixnet/build_setup.sh new file mode 100755 index 000000000..81af9d16f --- /dev/null +++ b/simulations/mixnet/build_setup.sh @@ -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 diff --git a/simulations/mixnet/config2.toml b/simulations/mixnet/config2.toml index c40e41103..3acd2bf8a 100644 --- a/simulations/mixnet/config2.toml +++ b/simulations/mixnet/config2.toml @@ -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 diff --git a/simulations/mixnet/config3.toml b/simulations/mixnet/config3.toml index 80c19b34b..bd8e7c4e9 100644 --- a/simulations/mixnet/config3.toml +++ b/simulations/mixnet/config3.toml @@ -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 diff --git a/simulations/mixnet/config4.toml b/simulations/mixnet/config4.toml index ed5b2dad0..f174250d5 100644 --- a/simulations/mixnet/config4.toml +++ b/simulations/mixnet/config4.toml @@ -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 diff --git a/simulations/mixnet/run_chat_mix.sh b/simulations/mixnet/run_chat_mix.sh index f711c055e..ef0575375 100755 --- a/simulations/mixnet/run_chat_mix.sh +++ b/simulations/mixnet/run_chat_mix.sh @@ -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" diff --git a/simulations/mixnet/run_chat_mix1.sh b/simulations/mixnet/run_chat_mix1.sh index 7323bb3a9..5961fce45 100755 --- a/simulations/mixnet/run_chat_mix1.sh +++ b/simulations/mixnet/run_chat_mix1.sh @@ -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" diff --git a/simulations/mixnet/run_mix_node.sh b/simulations/mixnet/run_mix_node.sh index 2b293540c..5d9ff70d8 100755 --- a/simulations/mixnet/run_mix_node.sh +++ b/simulations/mixnet/run_mix_node.sh @@ -1 +1,2 @@ ../../build/wakunode2 --config-file="config.toml" 2>&1 | tee mix_node.log + diff --git a/simulations/mixnet/setup_credentials.nim b/simulations/mixnet/setup_credentials.nim new file mode 100644 index 000000000..77c796354 --- /dev/null +++ b/simulations/mixnet/setup_credentials.nim @@ -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() diff --git a/tests/test_helpers.nim b/tests/test_helpers.nim index a4bb69fbc..a0442d732 100644 --- a/tests/test_helpers.nim +++ b/tests/test_helpers.nim @@ -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 = diff --git a/tests/test_peer_manager.nim b/tests/test_peer_manager.nim index d4c2af5b5..487d6c051 100644 --- a/tests/test_peer_manager.nim +++ b/tests/test_peer_manager.nim @@ -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, ) diff --git a/tests/test_waku_switch.nim b/tests/test_waku_switch.nim index c3f635c17..86faae83d 100644 --- a/tests/test_waku_switch.nim +++ b/tests/test_waku_switch.nim @@ -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) diff --git a/tests/waku_filter_v2/test_waku_client.nim b/tests/waku_filter_v2/test_waku_client.nim index c5bd3c558..8742d3898 100644 --- a/tests/waku_filter_v2/test_waku_client.nim +++ b/tests/waku_filter_v2/test_waku_client.nim @@ -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() diff --git a/tests/waku_filter_v2/test_waku_filter_dos_protection.nim b/tests/waku_filter_v2/test_waku_filter_dos_protection.nim index c55a9b4cd..d0f89ffd1 100644 --- a/tests/waku_filter_v2/test_waku_filter_dos_protection.nim +++ b/tests/waku_filter_v2/test_waku_filter_dos_protection.nim @@ -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)) ) diff --git a/tests/wakunode_rest/test_rest_store.nim b/tests/wakunode_rest/test_rest_store.nim index 0f4ed10ea..3cd114893 100644 --- a/tests/wakunode_rest/test_rest_store.nim +++ b/tests/wakunode_rest/test_rest_store.nim @@ -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)