diff --git a/logos_delivery/api/events/kernel_events.nim b/logos_delivery/api/events/kernel_events.nim index d8dd46627..5b9b0af16 100644 --- a/logos_delivery/api/events/kernel_events.nim +++ b/logos_delivery/api/events/kernel_events.nim @@ -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: diff --git a/logos_delivery/messaging/delivery_service/recv_service/recv_service.nim b/logos_delivery/messaging/delivery_service/recv_service/recv_service.nim index d2c369d46..63888c321 100644 --- a/logos_delivery/messaging/delivery_service/recv_service/recv_service.nim +++ b/logos_delivery/messaging/delivery_service/recv_service/recv_service.nim @@ -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) diff --git a/logos_delivery/waku/node/subscription_manager.nim b/logos_delivery/waku/node/subscription_manager.nim index 2c18a7b5e..e88aa5f96 100644 --- a/logos_delivery/waku/node/subscription_manager.nim +++ b/logos_delivery/waku/node/subscription_manager.nim @@ -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 diff --git a/logos_delivery/waku/node/waku_node.nim b/logos_delivery/waku/node/waku_node.nim index db90ab47c..adf4bed49 100644 --- a/logos_delivery/waku/node/waku_node.nim +++ b/logos_delivery/waku/node/waku_node.nim @@ -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()