diff --git a/logos_delivery/waku/rest_api/endpoint/relay/handlers.nim b/logos_delivery/waku/rest_api/endpoint/relay/handlers.nim index 6b1fdcc69..43da58212 100644 --- a/logos_delivery/waku/rest_api/endpoint/relay/handlers.nim +++ b/logos_delivery/waku/rest_api/endpoint/relay/handlers.nim @@ -179,8 +179,27 @@ proc installRelayApiHandlers*( "Failed to publish: error appending RLN proof to message: " & $error ) - (await node.wakuRelay.validateMessage(pubsubTopic, message)).isOkOr: - return RestApiResponse.badRequest("Failed to publish: " & error) + let firstValidateResult = await node.wakuRelay.validateMessage(pubsubTopic, message) + if firstValidateResult.isErr(): + if node.rln.isNil() or not firstValidateResult.error.contains( + RlnValidatorErrorMsg + ): + return + RestApiResponse.badRequest("Failed to publish: " & firstValidateResult.error) + # Stale RLN merkle root; force-refresh the cached path and retry validation once + info "relay publish rejected as RLN-invalid; refreshing merkle proof and retrying once" + message.proof = ( + await node.rln.generateRLNProof( + message.toRLNSignal(), + float64(getTime().toUnix()), + forceMerkleProofRefresh = true, + ) + ).valueOr: + return RestApiResponse.internalServerError( + "Failed to publish: error appending RLN proof on retry: " & $error + ) + (await node.wakuRelay.validateMessage(pubsubTopic, message)).isOkOr: + return RestApiResponse.badRequest("Failed to publish: " & error) # Log for message tracking purposes logMessageInfo(node.wakuRelay, "rest", pubsubTopic, "none", message, onRecv = true) @@ -308,8 +327,27 @@ proc installRelayApiHandlers*( "Failed to publish: error appending RLN proof to message: " & error ) - (await node.wakuRelay.validateMessage(pubsubTopic, message)).isOkOr: - return RestApiResponse.badRequest("Failed to publish: " & error) + let firstValidateResult = await node.wakuRelay.validateMessage(pubsubTopic, message) + if firstValidateResult.isErr(): + if node.rln.isNil() or not firstValidateResult.error.contains( + RlnValidatorErrorMsg + ): + return + RestApiResponse.badRequest("Failed to publish: " & firstValidateResult.error) + # Stale RLN merkle root; force-refresh the cached path and retry validation once + info "relay publish rejected as RLN-invalid; refreshing merkle proof and retrying once" + message.proof = ( + await node.rln.generateRLNProof( + message.toRLNSignal(), + float64(getTime().toUnix()), + forceMerkleProofRefresh = true, + ) + ).valueOr: + return RestApiResponse.internalServerError( + "Failed to publish: error appending RLN proof on retry: " & error + ) + (await node.wakuRelay.validateMessage(pubsubTopic, message)).isOkOr: + return RestApiResponse.badRequest("Failed to publish: " & error) # Log for message tracking purposes logMessageInfo(node.wakuRelay, "rest", pubsubTopic, "none", message, onRecv = true) diff --git a/logos_delivery/waku/rln/rln.nim b/logos_delivery/waku/rln/rln.nim index dab5ec3ef..c373f348f 100644 --- a/logos_delivery/waku/rln/rln.nim +++ b/logos_delivery/waku/rln/rln.nim @@ -233,10 +233,23 @@ proc mount( proc( msg: WakuMessage, senderEpochTime: float64 ): Future[Result[RequestGenerateRlnProof, string]] {.async.} = - let proof = (await rln.generateRLNProof(msg.toRLNSignal(), senderEpochTime)).valueOr: + let proofBytes = (await rln.generateRLNProof(msg.toRLNSignal(), senderEpochTime)).valueOr: return err("Could not create RLN proof: " & error) - return ok(RequestGenerateRlnProof(proof: proof)), + let rlnProof = RateLimitProof.init(proofBytes).valueOr: + return err("Could not decode RLN proof for root check: " & $error) + if await rln.groupManager.validateRoot(rlnProof.merkleRoot): + return ok(RequestGenerateRlnProof(proof: proofBytes)) + + # Cached merkle proof path root has slid out of the valid window; force-refresh and regenerate + info "RLN broker provider: stale merkle root detected; force-refreshing merkle path" + let retryProof = ( + await rln.generateRLNProof( + msg.toRLNSignal(), senderEpochTime, forceMerkleProofRefresh = true + ) + ).valueOr: + return err("Could not create RLN proof on retry: " & error) + return ok(RequestGenerateRlnProof(proof: retryProof)), ).isOkOr: return err("Proof generator provider cannot be set: " & $error) diff --git a/tests/waku_rln_relay/test_wakunode_rln_relay.nim b/tests/waku_rln_relay/test_wakunode_rln_relay.nim index 7954317e7..4c02b4cbd 100644 --- a/tests/waku_rln_relay/test_wakunode_rln_relay.nim +++ b/tests/waku_rln_relay/test_wakunode_rln_relay.nim @@ -11,7 +11,8 @@ import brokers/broker_context import - logos_delivery/waku/[waku_core, waku_node, rln], + logos_delivery/waku/[waku_core, waku_node, rln, rln/protocol_types], + logos_delivery/waku/requests/rln_requests, ../testlib/[wakucore, futures, wakunode, testutils], ./utils_onchain, ./rln/waku_rln_relay_utils @@ -751,3 +752,55 @@ procSuite "WakuNode - RLN relay": # Cleanup waitFor allFutures(node1.stop(), node2.stop()) + + asyncTest "broker proof provider retries with force-refresh when initial proof has stale root": + ## Exercises the reactive mechanism added to RequestGenerateRlnProof.setProvider + ## in rln.nim: when the cached Merkle proof path produces a proof with a root + ## that validateRoot rejects, the provider must force-refresh the path and + ## return a proof whose root is in the valid-roots window. + lockNewGlobalBrokerContext: + let nodeKey = generateSecp256k1Key() + let node = newTestWakuNode(nodeKey) + (await node.mountRelay()).isOkOr: + assert false, "Failed to mount relay" + + let wakuRlnConfig = getWakuRlnConfig( + manager = manager, + index = MembershipIndex(1), + epochSizeSec = 600, + userMessageLimit = 20, + ) + await node.setRlnValidator(wakuRlnConfig) + await node.start() + + let rlnManager = cast[OnchainGroupManager](node.rln.groupManager) + let idCredentials = generateCredentials() + (waitFor rlnManager.register(idCredentials, UserMessageLimit(20))).isOkOr: + assert false, "Failed to register: " & error + + let rootUpdated = waitFor rlnManager.updateRoots() + info "Updated root", rootUpdated + + let proofRes = waitFor rlnManager.fetchMerkleProofElements() + assert proofRes.isOk(), "failed to fetch merkle proof: " & proofRes.error + let goodCache = proofRes.get() + rlnManager.merkleProofCache = goodCache + + # Corrupt the cache so the first generateRLNProof inside the provider + # produces a proof with a Merkle root that is not in the valid-roots window. + # The provider must detect this via validateRoot, force-refresh, and retry. + rlnManager.merkleProofCache = newSeq[byte](goodCache.len) + + let msg = fakeWakuMessage() + let proofResult = + await RequestGenerateRlnProof.request(node.rln.brokerCtx, msg, epochTime()) + + check proofResult.isOk() + # The force-refresh inside the provider restored the correct path + check rlnManager.merkleProofCache == goodCache + # The returned proof carries a Merkle root that is in the valid-roots window + let rlnProof = RateLimitProof.init(proofResult.get().proof).get() + let rootValid = await node.rln.groupManager.validateRoot(rlnProof.merkleRoot) + check rootValid + + await node.stop() diff --git a/tests/wakunode_rest/test_rest_relay.nim b/tests/wakunode_rest/test_rest_relay.nim index ccf59b4fd..95f8b8370 100644 --- a/tests/wakunode_rest/test_rest_relay.nim +++ b/tests/wakunode_rest/test_rest_relay.nim @@ -793,3 +793,161 @@ suite "Waku v2 Rest API - Relay": await restServer.stop() await restServer.closeWait() await node.stop() + + asyncTest "Stale RLN proof triggers force-refresh and retry - POST /relay/v1/messages/{topic}": + ## When the cached Merkle proof path is stale the handler generates a proof + ## whose root the local RLN validator rejects. The handler must detect the + ## RlnValidatorErrorMsg, force-refresh the cached path, regenerate the proof + ## with the correct root, and succeed — returning 200 OK. + let node = testWakuNode() + (await node.mountRelay()).isOkOr: + assert false, "Failed to mount relay" + let wakuRlnConfig = getWakuRlnConfig( + manager = manager, + index = MembershipIndex(1), + epochSizeSec = 600, + userMessageLimit = 20, + ) + await node.setRlnValidator(wakuRlnConfig) + await node.start() + + let manager = cast[OnchainGroupManager](node.rln.groupManager) + let idCredentials = generateCredentials() + (waitFor manager.register(idCredentials, UserMessageLimit(20))).isOkOr: + assert false, "Failed to register: " & getCurrentExceptionMsg() + + let rootUpdated = waitFor manager.updateRoots() + info "Updated root", rootUpdated + + let proofRes = waitFor manager.fetchMerkleProofElements() + assert proofRes.isOk(), "failed to fetch merkle proof: " & proofRes.error + let goodCache = proofRes.get() + manager.merkleProofCache = goodCache + + # Corrupt the cache with zeros so the first generateRLNProof call produces a + # proof with a Merkle root that is not in the valid-roots window. + # validateMessage will return RlnValidatorErrorMsg, triggering the retry. + manager.merkleProofCache = newSeq[byte](goodCache.len) + + var restPort = Port(0) + let restAddress = parseIpAddress("0.0.0.0") + let restServer = WakuRestServerRef.init(restAddress, restPort).tryGet() + restPort = restServer.httpServer.address.port + let cache = MessageCache.init() + installRelayApiHandlers(restServer.router, node, cache) + restServer.start() + let client = newRestHttpClient(initTAddress(restAddress, restPort)) + + let simpleHandler = proc( + topic: PubsubTopic, msg: WakuMessage + ): Future[void] {.async, gcsafe.} = + await sleepAsync(0.milliseconds) + + node.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: + assert false, "Failed to subscribe to pubsub topic" + + let response = await client.relayPostMessagesV1( + DefaultPubsubTopic, + RelayWakuMessage( + payload: base64.encode("TEST-PAYLOAD"), + contentTopic: some(DefaultContentTopic), + timestamp: some(now()), + ), + ) + + # Handler force-refreshed the path and retried successfully + check: + response.status == 200 + $response.contentType == $MIMETYPE_TEXT + response.data == "OK" + manager.merkleProofCache == goodCache # force-refresh restored the correct path + + await restServer.stop() + await restServer.closeWait() + await node.stop() + + asyncTest "Stale RLN proof triggers force-refresh and retry - POST /relay/v1/auto/messages/{topic}": + ## Same reactive retry as the static-sharding handler, exercised via the + ## auto-sharding endpoint. A relay-only mesh node is connected so that + ## node.publish() has a gossipsub peer and can return success. + + # Relay-only mesh node — no RLN needed, just provides a gossipsub peer. + let meshNode = testWakuNode() + (await meshNode.mountRelay()).isOkOr: + assert false, "Failed to mount relay on mesh node" + require meshNode.mountAutoSharding(1, 8).isOk + await meshNode.start() + let meshHandler = proc( + topic: PubsubTopic, msg: WakuMessage + ): Future[void] {.async, gcsafe.} = + discard + meshNode.subscribe((kind: ContentSub, topic: DefaultContentTopic), meshHandler).isOkOr: + assert false, "Failed to subscribe mesh node" + + var node: WakuNode + lockNewGlobalBrokerContext: + node = testWakuNode() + (await node.mountRelay()).isOkOr: + assert false, "Failed to mount relay" + require node.mountAutoSharding(1, 8).isOk + + let wakuRlnConfig = getWakuRlnConfig( + manager = manager, + index = MembershipIndex(1), + epochSizeSec = 600, + userMessageLimit = 20, + ) + await node.setRlnValidator(wakuRlnConfig) + await node.start() + await node.connectToNodes(@[meshNode.peerInfo.toRemotePeerInfo()]) + + let manager = cast[OnchainGroupManager](node.rln.groupManager) + let idCredentials = generateCredentials() + (waitFor manager.register(idCredentials, UserMessageLimit(20))).isOkOr: + assert false, "Failed to register: " & getCurrentExceptionMsg() + + let rootUpdated = waitFor manager.updateRoots() + info "Updated root", rootUpdated + + let proofRes = waitFor manager.fetchMerkleProofElements() + assert proofRes.isOk(), "failed to fetch merkle proof: " & proofRes.error + let goodCache = proofRes.get() + manager.merkleProofCache = goodCache + + # Corrupt the cache to produce a proof with a bad Merkle root + manager.merkleProofCache = newSeq[byte](goodCache.len) + + var restPort = Port(0) + let restAddress = parseIpAddress("0.0.0.0") + let restServer = WakuRestServerRef.init(restAddress, restPort).tryGet() + restPort = restServer.httpServer.address.port + let cache = MessageCache.init() + installRelayApiHandlers(restServer.router, node, cache) + restServer.start() + let client = newRestHttpClient(initTAddress(restAddress, restPort)) + + let simpleHandler = proc( + topic: PubsubTopic, msg: WakuMessage + ): Future[void] {.async, gcsafe.} = + await sleepAsync(0.milliseconds) + + node.subscribe((kind: ContentSub, topic: DefaultContentTopic), simpleHandler).isOkOr: + assert false, "Failed to subscribe to content topic" + + let response = await client.relayPostAutoMessagesV1( + RelayWakuMessage( + payload: base64.encode("TEST-PAYLOAD"), + contentTopic: some(DefaultContentTopic), + timestamp: some(now()), + ) + ) + + check: + response.status == 200 + $response.contentType == $MIMETYPE_TEXT + response.data == "OK" + manager.merkleProofCache == goodCache + + await restServer.stop() + await restServer.closeWait() + await allFutures(node.stop(), meshNode.stop())