mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-22 04:29:42 +00:00
Extend reactive Merkle proof refresh to REST relay and broker paths
With PathCheckMinInterval removed, a non-empty merkleProofCache is trusted forever unless force=true is passed. Previously only the two lightpush retry paths ever passed force=true; the REST relay publish handlers (static and auto-sharding) and the broker proof provider had no rejection feedback loop, so a sliding root window would permanently reject their messages until restart. - REST relay handlers: run validateMessage first; if it returns RlnValidatorErrorMsg, force-refresh the cached path and retry validation once before publishing — mirroring the lightpush reactive pattern. - Broker provider (rln.nim): decode the generated proof bytes, call validateRoot on the embedded Merkle root, and force-refresh + regenerate when the root is outside the acceptance window — since the broker has no external rejection signal to react to. Tests added for all three new paths, using the corrupted-cache technique (all-zero merkleProofCache → garbage root → rejection → retry). Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
e7afdce63e
commit
f02cb1f881
@ -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)
|
||||
|
||||
@ -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)
|
||||
|
||||
|
||||
@ -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()
|
||||
|
||||
@ -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())
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user