mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-22 12:39:30 +00:00
Break up legacyLightpushPublish into single-purpose helpers
legacyLightpushPublish had grown to ~100 lines with three concerns
interleaved: choosing between client/self-request publish, resolving the
effective pubsub topic, and running the RLN refresh-retry with timeout.
Lift each into a private top-level proc so the main proc reads top-to-
bottom in ~35 lines: precondition → prepare message + proof → resolve
topic → publish once → RLN-refresh retry if applicable.
The new helpers are:
- internalLegacyLightpushPublish — client vs self-request dispatch
(was the inline closure)
- resolveLegacyPubsubTopic — explicit param or autosharding
- runRlnRefreshRetry — force-refresh proof + retry with
RlnRefreshRetryTimeout bound; returns caller's fallback on timeout
Nonce now appears in a top-level signature, so import it directly from
rln/nonce_manager rather than re-exporting through rln/proof.nim (the
latter would leak nonce_manager's bulk chronos/times exports and clash
with waku_rendezvous overload resolution).
Pure refactor; behavior unchanged. Verified against
tests/node/test_wakunode_legacy_lightpush.nim (9/9) and
tests/node/test_wakunode_lightpush.nim (11/11).
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
parent
af94b3c3cf
commit
61bcdf6c6e
@ -30,7 +30,8 @@ import
|
||||
../../waku_lightpush as lightpush_protocol,
|
||||
../peer_manager,
|
||||
../../common/rate_limit/setting,
|
||||
../../rln
|
||||
../../rln,
|
||||
../../rln/nonce_manager
|
||||
|
||||
logScope:
|
||||
topics = "waku node lightpush api"
|
||||
@ -68,6 +69,80 @@ proc mountLegacyLightPushClient*(node: WakuNode) =
|
||||
node.wakuLegacyLightpushClient =
|
||||
WakuLegacyLightPushClient.new(node.peerManager, node.rng)
|
||||
|
||||
proc internalLegacyLightpushPublish(
|
||||
node: WakuNode,
|
||||
pubsubTopic: PubsubTopic,
|
||||
message: WakuMessage,
|
||||
peer: RemotePeerInfo,
|
||||
): Future[legacy_lightpush_protocol.WakuLightPushResult[string]] {.async, gcsafe.} =
|
||||
## Dispatches to the legacy lightpush client if mounted, otherwise to the
|
||||
## self-hosted server. Callers guarantee at least one is mounted.
|
||||
let msgHash = pubsubTopic.computeMessageHash(message).to0xHex()
|
||||
if not node.wakuLegacyLightpushClient.isNil():
|
||||
notice "publishing message with legacy lightpush",
|
||||
pubsubTopic = pubsubTopic,
|
||||
contentTopic = message.contentTopic,
|
||||
target_peer_id = peer.peerId,
|
||||
msg_hash = msgHash
|
||||
return await node.wakuLegacyLightpushClient.publish(pubsubTopic, message, peer)
|
||||
|
||||
notice "publishing message with self hosted legacy lightpush",
|
||||
pubsubTopic = pubsubTopic,
|
||||
contentTopic = message.contentTopic,
|
||||
target_peer_id = peer.peerId,
|
||||
msg_hash = msgHash
|
||||
return await node.wakuLegacyLightPush.handleSelfLightPushRequest(pubsubTopic, message)
|
||||
|
||||
proc resolveLegacyPubsubTopic(
|
||||
node: WakuNode, pubsubTopic: Option[PubsubTopic], contentTopic: ContentTopic
|
||||
): Result[PubsubTopic, string] =
|
||||
## Returns the explicit pubsub topic if provided, otherwise derives it from
|
||||
## `contentTopic` via autosharding.
|
||||
if pubsubTopic.isSome():
|
||||
return ok(pubsubTopic.get())
|
||||
if node.wakuAutoSharding.isNone():
|
||||
return err("Pubsub topic must be specified when static sharding is enabled")
|
||||
let parsedTopic = NsContentTopic.parse(contentTopic).valueOr:
|
||||
return err("Invalid content-topic: " & $error)
|
||||
let shard = node.wakuAutoSharding.get().getShard(parsedTopic).valueOr:
|
||||
return err("Autosharding error: " & error)
|
||||
return ok($shard)
|
||||
|
||||
proc runRlnRefreshRetry(
|
||||
node: WakuNode,
|
||||
rln: Option[Rln],
|
||||
msgWithProof: WakuMessage,
|
||||
drawnMessageId: Option[Nonce],
|
||||
pubsubForPublish: PubsubTopic,
|
||||
peer: RemotePeerInfo,
|
||||
fallback: legacy_lightpush_protocol.WakuLightPushResult[string],
|
||||
): Future[legacy_lightpush_protocol.WakuLightPushResult[string]] {.async, gcsafe.} =
|
||||
## Force-refreshes the RLN merkle proof path and retries the publish once,
|
||||
## bounded by RlnRefreshRetryTimeout. Returns `fallback` on timeout so a
|
||||
## hanging RPC/libp2p round-trip cannot stall the caller indefinitely.
|
||||
info "legacy lightpush send rejected as RLN-invalid; " &
|
||||
"refreshing merkle proof and retrying once"
|
||||
|
||||
proc runRetry(): Future[legacy_lightpush_protocol.WakuLightPushResult[string]] {.
|
||||
async, gcsafe
|
||||
.} =
|
||||
let (retryMsg, _) = (
|
||||
await checkAndGenerateRLNProof(
|
||||
rln,
|
||||
msgWithProof,
|
||||
forceMerkleProofRefresh = true,
|
||||
reuseMessageId = drawnMessageId,
|
||||
)
|
||||
).valueOr:
|
||||
return err("failed call checkAndGenerateRLNProof from lightpush retry: " & error)
|
||||
return await internalLegacyLightpushPublish(node, pubsubForPublish, retryMsg, peer)
|
||||
|
||||
let retryFut = runRetry()
|
||||
if not (await retryFut.withTimeout(RlnRefreshRetryTimeout)):
|
||||
warn "legacy lightpush RLN-refresh retry timed out; returning original error"
|
||||
return fallback
|
||||
return retryFut.read()
|
||||
|
||||
proc legacyLightpushPublish*(
|
||||
node: WakuNode,
|
||||
pubsubTopic: Option[PubsubTopic],
|
||||
@ -82,9 +157,9 @@ proc legacyLightpushPublish*(
|
||||
error "failed to publish message as legacy lightpush not available"
|
||||
return err("Waku lightpush not available")
|
||||
|
||||
# toRLNSignal includes the timestamp in the proof input, so the timestamp
|
||||
# must be fixed before proof generation. The downstream ensureTimestampSet
|
||||
# in the client publish becomes an idempotent no-op safety net.
|
||||
# toRLNSignal includes the timestamp in the proof input, so it must be fixed
|
||||
# before proof generation. The downstream ensureTimestampSet in the client
|
||||
# publish becomes an idempotent no-op safety net.
|
||||
let message = ensureTimestampSet(message)
|
||||
|
||||
let rln =
|
||||
@ -95,46 +170,15 @@ proc legacyLightpushPublish*(
|
||||
let (msgWithProof, drawnMessageId) = (await checkAndGenerateRLNProof(rln, message)).valueOr:
|
||||
return err("failed call checkAndGenerateRLNProof from lightpush: " & error)
|
||||
|
||||
let internalPublish = proc(
|
||||
node: WakuNode,
|
||||
pubsubTopic: PubsubTopic,
|
||||
message: WakuMessage,
|
||||
peer: RemotePeerInfo,
|
||||
): Future[legacy_lightpush_protocol.WakuLightPushResult[string]] {.async, gcsafe.} =
|
||||
let msgHash = pubsubTopic.computeMessageHash(message).to0xHex()
|
||||
if not node.wakuLegacyLightpushClient.isNil():
|
||||
notice "publishing message with legacy lightpush",
|
||||
pubsubTopic = pubsubTopic,
|
||||
contentTopic = message.contentTopic,
|
||||
target_peer_id = peer.peerId,
|
||||
msg_hash = msgHash
|
||||
return await node.wakuLegacyLightpushClient.publish(pubsubTopic, message, peer)
|
||||
|
||||
if not node.wakuLegacyLightPush.isNil():
|
||||
notice "publishing message with self hosted legacy lightpush",
|
||||
pubsubTopic = pubsubTopic,
|
||||
contentTopic = message.contentTopic,
|
||||
target_peer_id = peer.peerId,
|
||||
msg_hash = msgHash
|
||||
return
|
||||
await node.wakuLegacyLightPush.handleSelfLightPushRequest(pubsubTopic, message)
|
||||
try:
|
||||
# Resolve the effective pubsub topic first (explicit param, else derive
|
||||
# from autosharding), then run the send. We need a single resolved topic
|
||||
# so the retry path below can reuse it without repeating this branch.
|
||||
var pubsubForPublish: PubsubTopic
|
||||
if pubsubTopic.isSome():
|
||||
pubsubForPublish = pubsubTopic.get()
|
||||
else:
|
||||
if node.wakuAutoSharding.isNone():
|
||||
return err("Pubsub topic must be specified when static sharding is enabled")
|
||||
let parsedTopic = NsContentTopic.parse(message.contentTopic).valueOr:
|
||||
return err("Invalid content-topic: " & $error)
|
||||
let shard = node.wakuAutoSharding.get().getShard(parsedTopic).valueOr:
|
||||
return err("Autosharding error: " & error)
|
||||
pubsubForPublish = $shard
|
||||
let pubsubForPublish = resolveLegacyPubsubTopic(
|
||||
node, pubsubTopic, message.contentTopic
|
||||
).valueOr:
|
||||
return err(error)
|
||||
|
||||
let firstResult = await internalPublish(node, pubsubForPublish, msgWithProof, peer)
|
||||
let firstResult = await internalLegacyLightpushPublish(
|
||||
node, pubsubForPublish, msgWithProof, peer
|
||||
)
|
||||
|
||||
# A publish error mentioning RLN can indicate a stale merkle proof path;
|
||||
# refresh it and retry the publish once. Legacy lightpush has no status
|
||||
@ -143,29 +187,9 @@ proc legacyLightpushPublish*(
|
||||
not firstResult.error.contains(RlnValidatorErrorMsg):
|
||||
return firstResult
|
||||
|
||||
info "legacy lightpush send rejected as RLN-invalid; " &
|
||||
"refreshing merkle proof and retrying once"
|
||||
|
||||
proc runRetry(): Future[legacy_lightpush_protocol.WakuLightPushResult[string]] {.
|
||||
async, gcsafe
|
||||
.} =
|
||||
let (retryMsg, _) = (
|
||||
await checkAndGenerateRLNProof(
|
||||
rln,
|
||||
msgWithProof,
|
||||
forceMerkleProofRefresh = true,
|
||||
reuseMessageId = drawnMessageId,
|
||||
)
|
||||
).valueOr:
|
||||
return
|
||||
err("failed call checkAndGenerateRLNProof from lightpush retry: " & error)
|
||||
return await internalPublish(node, pubsubForPublish, retryMsg, peer)
|
||||
|
||||
let retryFut = runRetry()
|
||||
if not (await retryFut.withTimeout(RlnRefreshRetryTimeout)):
|
||||
warn "legacy lightpush RLN-refresh retry timed out; returning original error"
|
||||
return firstResult
|
||||
return retryFut.read()
|
||||
return await runRlnRefreshRetry(
|
||||
node, rln, msgWithProof, drawnMessageId, pubsubForPublish, peer, firstResult
|
||||
)
|
||||
except CatchableError:
|
||||
return err(getCurrentExceptionMsg())
|
||||
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user