mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-23 05:00:21 +00:00
refactor(node): thread WakuEnvelope through dispatch; MessageSeenEvent carries it
subscription_manager relay handlers take the envelope (one async capture instead of topic+msg), reusing envelope.hash for store-sync ingress. MessageSeenEvent now carries the WakuEnvelope; the filter-client push emit builds one, and the recv_service listener reuses envelope.hash (optional precomputed-hash param) to skip recomputation. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
8a2172a2eb
commit
1913a36007
@ -7,10 +7,11 @@ import logos_delivery/waku/waku_core/message
|
||||
export event_broker, pubsub_topic, message
|
||||
|
||||
EventBroker:
|
||||
# Internal event emitted when a message arrives from the network via any protocol
|
||||
# Internal event emitted when a message arrives from the network via any protocol.
|
||||
# Carries the WakuEnvelope so listeners reuse the precomputed hash instead of
|
||||
# recomputing it.
|
||||
type MessageSeenEvent* = object
|
||||
topic*: PubsubTopic
|
||||
message*: WakuMessage
|
||||
envelope*: WakuEnvelope
|
||||
|
||||
# Emitted by the health monitor when overall node connectivity changes.
|
||||
EventBroker:
|
||||
|
||||
@ -73,18 +73,24 @@ proc getMissingMsgsFromStore(
|
||||
)
|
||||
|
||||
proc processIncomingMessage(
|
||||
self: RecvService, pubsubTopic: string, message: WakuMessage
|
||||
self: RecvService,
|
||||
pubsubTopic: string,
|
||||
message: WakuMessage,
|
||||
precomputedHash = Opt.none(WakuMessageHash),
|
||||
): bool =
|
||||
## Return false if the incoming message is from a non-subscribed topic,
|
||||
## or if the message is a duplicate (recently-seen). Otherwise, save it as
|
||||
## recently-seen, emit a MessageReceivedEvent, and return true.
|
||||
## `precomputedHash` lets callers that already have the hash (e.g. a
|
||||
## MessageSeenEvent envelope) skip recomputation.
|
||||
|
||||
if not self.waku.isContentSubscribed(pubsubTopic, message.contentTopic):
|
||||
trace "skipping message as I am not subscribed",
|
||||
shard = pubsubTopic, contentTopic = message.contentTopic
|
||||
return false
|
||||
|
||||
let msgHash = computeMessageHash(pubsubTopic, message)
|
||||
let msgHash = precomputedHash.valueOr:
|
||||
computeMessageHash(pubsubTopic, message)
|
||||
if self.recentReceivedMsgs.anyIt(it.msgHash == msgHash):
|
||||
trace "skipping duplicate message",
|
||||
shard = pubsubTopic,
|
||||
@ -188,7 +194,9 @@ proc startRecvService*(self: RecvService) =
|
||||
self.seenMsgListener = MessageSeenEvent.listen(
|
||||
self.brokerCtx,
|
||||
proc(event: MessageSeenEvent) {.async: (raises: []).} =
|
||||
discard self.processIncomingMessage(event.topic, event.message),
|
||||
discard self.processIncomingMessage(
|
||||
event.envelope.pubsubTopic, event.envelope.msg, Opt.some(event.envelope.hash)
|
||||
),
|
||||
).valueOr:
|
||||
error "Failed to set MessageSeenEvent listener", error = error
|
||||
quit(QuitFailure)
|
||||
|
||||
@ -42,44 +42,48 @@ proc registerRelayHandler(
|
||||
if alreadySubscribed:
|
||||
return false
|
||||
|
||||
proc traceHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
let msgSizeKB = msg.payload.len / 1000
|
||||
proc traceHandler(envelope: WakuEnvelope) {.async, gcsafe.} =
|
||||
let msgSizeKB = envelope.msg.payload.len / 1000
|
||||
|
||||
waku_node_messages.inc(labelValues = ["relay"])
|
||||
waku_histogram_message_size.observe(msgSizeKB)
|
||||
|
||||
proc filterHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
proc filterHandler(envelope: WakuEnvelope) {.async, gcsafe.} =
|
||||
if node.wakuFilter.isNil():
|
||||
return
|
||||
|
||||
await node.wakuFilter.handleMessage(topic, msg)
|
||||
await node.wakuFilter.handleMessage(envelope)
|
||||
|
||||
proc archiveHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
proc archiveHandler(envelope: WakuEnvelope) {.async, gcsafe.} =
|
||||
if node.wakuArchive.isNil():
|
||||
return
|
||||
|
||||
await node.wakuArchive.handleMessage(topic, msg)
|
||||
await node.wakuArchive.handleMessage(envelope)
|
||||
|
||||
proc syncHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
proc syncHandler(envelope: WakuEnvelope) {.async, gcsafe.} =
|
||||
if node.wakuStoreReconciliation.isNil():
|
||||
return
|
||||
|
||||
node.wakuStoreReconciliation.messageIngress(topic, msg)
|
||||
# Reuse the envelope's precomputed hash (no re-hash).
|
||||
node.wakuStoreReconciliation.messageIngress(
|
||||
envelope.hash, envelope.pubsubTopic, envelope.msg
|
||||
)
|
||||
|
||||
proc internalHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
MessageSeenEvent.emit(node.brokerCtx, topic, msg)
|
||||
proc internalHandler(envelope: WakuEnvelope) {.async, gcsafe.} =
|
||||
MessageSeenEvent.emit(node.brokerCtx, envelope)
|
||||
|
||||
let uniqueTopicHandler = proc(
|
||||
topic: PubsubTopic, msg: WakuMessage
|
||||
envelope: WakuEnvelope
|
||||
): Future[void] {.async, gcsafe.} =
|
||||
await traceHandler(topic, msg)
|
||||
await filterHandler(topic, msg)
|
||||
await archiveHandler(topic, msg)
|
||||
await syncHandler(topic, msg)
|
||||
await internalHandler(topic, msg)
|
||||
let topic = envelope.pubsubTopic
|
||||
await traceHandler(envelope)
|
||||
await filterHandler(envelope)
|
||||
await archiveHandler(envelope)
|
||||
await syncHandler(envelope)
|
||||
await internalHandler(envelope)
|
||||
|
||||
if node.legacyAppHandlers.hasKey(topic) and not node.legacyAppHandlers[topic].isNil():
|
||||
await node.legacyAppHandlers[topic](topic, msg)
|
||||
await node.legacyAppHandlers[topic](envelope)
|
||||
|
||||
node.wakuRelay.subscribe(shard, uniqueTopicHandler)
|
||||
return true
|
||||
|
||||
@ -633,7 +633,7 @@ proc start*(node: WakuNode) {.async.} =
|
||||
if not node.wakuFilterClient.isNil():
|
||||
node.wakuFilterClient.registerPushHandler(
|
||||
proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} =
|
||||
MessageSeenEvent.emit(node.brokerCtx, pubsubTopic, msg)
|
||||
MessageSeenEvent.emit(node.brokerCtx, WakuEnvelope.init(pubsubTopic, msg))
|
||||
)
|
||||
|
||||
node.startProvidersAndListeners()
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user