From 35635f98434ba6a6fa8a1b0a99ae7561be60f112 Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:44:13 +0200 Subject: [PATCH 01/10] feat(waku_core): add WakuEnvelope (msg + topic + hash, decode/hash-once carrier) Introduces the immutable-by-convention ref that flows through the internal relay dispatch so consumers reuse one decode and one hash. Includes unit tests asserting envelope.hash == computeMessageHash and ref (no-copy) semantics. Co-Authored-By: Claude Fable 5 --- logos_delivery/waku/waku_core/message.nim | 3 +- .../waku/waku_core/message/envelope.nim | 38 ++++++++++++++++ tests/all_tests_waku.nim | 1 + tests/waku_core/test_all.nim | 1 + tests/waku_core/test_message_envelope.nim | 44 +++++++++++++++++++ 5 files changed, 86 insertions(+), 1 deletion(-) create mode 100644 logos_delivery/waku/waku_core/message/envelope.nim create mode 100644 tests/waku_core/test_message_envelope.nim diff --git a/logos_delivery/waku/waku_core/message.nim b/logos_delivery/waku/waku_core/message.nim index 3f2d93d8a..c009aefa1 100644 --- a/logos_delivery/waku/waku_core/message.nim +++ b/logos_delivery/waku/waku_core/message.nim @@ -3,6 +3,7 @@ import ./message/default_values, ./message/codec, ./message/digest, + ./message/envelope, ./message/path_counters -export message, default_values, codec, digest, path_counters +export message, default_values, codec, digest, envelope, path_counters diff --git a/logos_delivery/waku/waku_core/message/envelope.nim b/logos_delivery/waku/waku_core/message/envelope.nim new file mode 100644 index 000000000..6aaf86e73 --- /dev/null +++ b/logos_delivery/waku/waku_core/message/envelope.nim @@ -0,0 +1,38 @@ +## Waku message envelope. +## +## Bundles a decoded `WakuMessage` with its `pubsubTopic` and the deterministic +## `WakuMessageHash`, computed **once** at construction. The envelope is the unit +## that flows through the internal relay dispatch (relay topic handler -> +## subscription_manager -> archive / filter / store-sync / app handlers) so that +## the same message is neither re-decoded nor re-hashed by each consumer. +## +## Like `WakuMessage`, a `WakuEnvelope` is a `ref object` and **immutable by +## convention**: construct it once after validation and never mutate it. Under +## `--mm:refc` passing it around is a pointer + refcount, not a deep copy. + +{.push raises: [].} + +import ../topics, ./message, ./digest + +type WakuEnvelope* = ref object + msg*: WakuMessage + pubsubTopic*: PubsubTopic + hash*: WakuMessageHash + +proc init*(T: type WakuEnvelope, pubsubTopic: PubsubTopic, msg: WakuMessage): T = + ## Builds an envelope, computing the message hash once (the single inbound-path + ## hash). `msg` is referenced, not copied. + WakuEnvelope( + msg: msg, pubsubTopic: pubsubTopic, hash: computeMessageHash(pubsubTopic, msg) + ) + +proc shortLog*(envelope: WakuEnvelope): string = + ## Compact chronicles representation: short hash + topic. + if envelope.isNil(): + return "nil" + "hash=" & envelope.hash.to0xHex() & " topic=" & envelope.pubsubTopic + +proc `$`*(envelope: WakuEnvelope): string = + shortLog(envelope) + +{.pop.} diff --git a/tests/all_tests_waku.nim b/tests/all_tests_waku.nim index f8ffc7b20..0e2b6ce4d 100644 --- a/tests/all_tests_waku.nim +++ b/tests/all_tests_waku.nim @@ -8,6 +8,7 @@ import ./waku_core/test_time, ./waku_core/test_message, ./waku_core/test_message_digest, + ./waku_core/test_message_envelope, ./waku_core/test_peers, ./waku_core/test_published_address diff --git a/tests/waku_core/test_all.nim b/tests/waku_core/test_all.nim index f7f4fad38..7ff2545dc 100644 --- a/tests/waku_core/test_all.nim +++ b/tests/waku_core/test_all.nim @@ -2,6 +2,7 @@ import ./test_message_digest, + ./test_message_envelope, ./test_namespaced_topics, ./test_peers, ./test_published_address, diff --git a/tests/waku_core/test_message_envelope.nim b/tests/waku_core/test_message_envelope.nim new file mode 100644 index 000000000..67b37d6fc --- /dev/null +++ b/tests/waku_core/test_message_envelope.nim @@ -0,0 +1,44 @@ +{.used.} + +import std/sequtils, stew/byteutils, testutils/unittests +import logos_delivery/waku/waku_core, ../testlib/wakucore + +suite "Waku Message - Envelope": + test "envelope init computes the same hash as computeMessageHash": + ## Given + let pubsubTopic = DefaultPubsubTopic + let message = fakeWakuMessage( + contentTopic = DefaultContentTopic, + payload = "\x01\x02\x03\x04TEST\x05\x06\x07\x08".toBytes(), + meta = newSeq[byte](), + ts = getNanosecondTime(1681964442), + ) + + ## When + let envelope = WakuEnvelope.init(pubsubTopic, message) + + ## Then + check: + envelope.hash == computeMessageHash(pubsubTopic, message) + envelope.pubsubTopic == pubsubTopic + envelope.msg == message + + test "envelope references the same message (no copy)": + let pubsubTopic = DefaultPubsubTopic + let message = fakeWakuMessage(payload = "abc".toBytes()) + let envelope = WakuEnvelope.init(pubsubTopic, message) + + ## The envelope holds the very same ref, not a clone. + check: + envelope.msg == message + # ref identity: mutating through one is visible through the other + cast[pointer](envelope.msg) == cast[pointer](message) + + test "different topics yield different hashes for the same message": + let message = fakeWakuMessage(payload = "same-payload".toBytes()) + let e1 = WakuEnvelope.init("/waku/2/rs/0/0", message) + let e2 = WakuEnvelope.init("/waku/2/rs/0/1", message) + + check: + e1.hash != e2.hash + e1.msg == e2.msg From 8a2172a2ebc4d62cd602bb61137ccec1c8bbed2d Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:44:20 +0200 Subject: [PATCH 02/10] refactor(waku_relay): WakuRelayHandler takes WakuEnvelope; decode/hash once Internal API break. The topic handler now builds the WakuEnvelope (the single inbound-path hash) and passes it to the registered handler. Deletes the unused onRecv observer decode (byte metrics need no decode) and gates the onValidated/ onSend observer decode+hash+per-shard gauges behind enabledLogLevel <= DEBUG so production INFO+ builds pay neither. Bench accessors unaffected (TopicHandler signature unchanged). Co-Authored-By: Claude Fable 5 --- logos_delivery/waku/waku_relay/protocol.nim | 120 ++++++++------------ 1 file changed, 47 insertions(+), 73 deletions(-) diff --git a/logos_delivery/waku/waku_relay/protocol.nim b/logos_delivery/waku/waku_relay/protocol.nim index 17c6e8e03..9e2592f9b 100644 --- a/logos_delivery/waku/waku_relay/protocol.nim +++ b/logos_delivery/waku/waku_relay/protocol.nim @@ -150,9 +150,8 @@ const GossipsubParameters = GossipSubParams.init( type WakuRelayResult*[T] = Result[T, string] - WakuRelayHandler* = proc(pubsubTopic: PubsubTopic, message: WakuMessage): Future[void] {. - gcsafe, raises: [Defect] - .} + WakuRelayHandler* = + proc(envelope: WakuEnvelope): Future[void] {.gcsafe, raises: [Defect].} WakuValidatorHandler* = proc( pubsubTopic: PubsubTopic, message: WakuMessage ): Future[ValidationResult] {.gcsafe, raises: [Defect].} @@ -253,51 +252,6 @@ proc logMessageInfo*( waku_relay_total_msg_bytes_per_shard.set(shardMetrics.sizeSum, labelValues = [topic]) proc initRelayObservers(w: WakuRelay) = - proc decodeRpcMessageInfo( - peer: PubSubPeer, msg: Message - ): Result[ - tuple[msgId: string, topic: string, wakuMessage: WakuMessage, msgSize: int], void - ] = - let msg_id = w.msgIdProvider(msg).valueOr: - warn "Error generating message id", - my_peer_id = w.switch.peerInfo.peerId, - from_peer_id = peer.peerId, - pubsub_topic = msg.topic, - error = $error - return err() - - let msg_id_short = shortLog(msg_id) - - let wakuMessage = WakuMessage.decode(msg.data).valueOr: - warn "Error decoding to Waku Message", - my_peer_id = w.switch.peerInfo.peerId, - msg_id = msg_id_short, - from_peer_id = peer.peerId, - pubsub_topic = msg.topic, - error = $error - return err() - - let msgSize = msg.data.len + msg.topic.len - return ok((msg_id_short, msg.topic, wakuMessage, msgSize)) - - proc updateMetrics( - peer: PubSubPeer, - pubsub_topic: string, - msg: WakuMessage, - msgSize: int, - onRecv: bool, - ) = - if onRecv: - waku_relay_network_bytes.inc( - msgSize.int64, labelValues = [pubsub_topic, "gross", "in"] - ) - else: - # sent traffic can only be "net" - # TODO: If we can measure unsuccessful sends would mean a possible distinction between gross/net - waku_relay_network_bytes.inc( - msgSize.int64, labelValues = [pubsub_topic, "net", "out"] - ) - proc onRecv(peer: PubSubPeer, msgs: var RPCMsg) = if msgs.control.isSome(): let ctrl = msgs.control.get() @@ -315,37 +269,53 @@ proc initRelayObservers(w: WakuRelay) = w.topicHealthUpdateEvent.fire() for msg in msgs.messages: - let (msg_id_short, topic, wakuMessage, msgSize) = decodeRpcMessageInfo(peer, msg).valueOr: - continue - # message receive log happens in onValidated observer as onRecv is called before checks - updateMetrics(peer, topic, wakuMessage, msgSize, onRecv = true) - discard + # gross incoming traffic; message size needs no proto decode (encoded + # buffer length + topic length). The receive log happens in onValidated. + waku_relay_network_bytes.inc( + (msg.data.len + msg.topic.len).int64, labelValues = [msg.topic, "gross", "in"] + ) proc onValidated(peer: PubSubPeer, msg: Message, msgId: MessageId) = - let msg_id_short = shortLog(msgId) - let wakuMessage = WakuMessage.decode(msg.data).valueOr: - warn "onValidated: failed decoding to Waku Message", - my_peer_id = w.switch.peerInfo.peerId, - msg_id = msg_id_short, - from_peer_id = peer.peerId, - pubsub_topic = msg.topic, - error = $error - return + # The per-message receive log + per-shard byte gauges require a full proto + # decode + hash. Gate them to DEBUG/TRACE builds so production (INFO+) pays + # neither. See docs/analysis/plan_phase3_wakuenvelope.md Step 2.4. + when enabledLogLevel <= LogLevel.DEBUG: + let msg_id_short = shortLog(msgId) + let wakuMessage = WakuMessage.decode(msg.data).valueOr: + warn "onValidated: failed decoding to Waku Message", + my_peer_id = w.switch.peerInfo.peerId, + msg_id = msg_id_short, + from_peer_id = peer.peerId, + pubsub_topic = msg.topic, + error = $error + return - logMessageInfo( - w, shortLog(peer.peerId), msg.topic, msg_id_short, wakuMessage, onRecv = true - ) + logMessageInfo( + w, shortLog(peer.peerId), msg.topic, msg_id_short, wakuMessage, onRecv = true + ) proc onSend(peer: PubSubPeer, msgs: var RPCMsg) = for msg in msgs.messages: - let (msg_id_short, topic, wakuMessage, msgSize) = decodeRpcMessageInfo(peer, msg).valueOr: - warn "onSend: failed decoding RPC info", - my_peer_id = w.switch.peerInfo.peerId, to_peer_id = peer.peerId - continue - logMessageInfo( - w, shortLog(peer.peerId), topic, msg_id_short, wakuMessage, onRecv = false + # net outgoing traffic; size needs no decode. + waku_relay_network_bytes.inc( + (msg.data.len + msg.topic.len).int64, labelValues = [msg.topic, "net", "out"] ) - updateMetrics(peer, topic, wakuMessage, msgSize, onRecv = false) + # The send log requires a decode + hash: gate to DEBUG/TRACE builds. + when enabledLogLevel <= LogLevel.DEBUG: + let msg_id = w.msgIdProvider(msg).valueOr: + continue + let wakuMessage = WakuMessage.decode(msg.data).valueOr: + warn "onSend: failed decoding to Waku Message", + my_peer_id = w.switch.peerInfo.peerId, to_peer_id = peer.peerId + continue + logMessageInfo( + w, + shortLog(peer.peerId), + msg.topic, + shortLog(msg_id), + wakuMessage, + onRecv = false, + ) let administrativeObserver = PubSubObserver(onRecv: onRecv, onSend: onSend, onValidated: onValidated) @@ -613,7 +583,11 @@ proc subscribe*(w: WakuRelay, pubsubTopic: PubsubTopic, handler: WakuRelayHandle data.len.int64 + pubsubTopic.len.int64, labelValues = [pubsubTopic, "net", "in"] ) - return handler(pubsubTopic, decMsg) + # Build the envelope here: the single inbound-path hash. It carries the + # decoded message + topic + hash through the whole dispatch chain so no + # downstream consumer re-decodes or re-hashes. + let envelope = WakuEnvelope.init(pubsubTopic, decMsg) + return handler(envelope) # Add the ordered validator to the topic # This assumes that if `w.validatorInserted.hasKey(pubSubTopic) is true`, it contains the ordered validator. From 1913a36007c6cf32c2dc54b03cdc528102abf169 Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:44:32 +0200 Subject: [PATCH 03/10] 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 --- logos_delivery/api/events/kernel_events.nim | 7 ++-- .../recv_service/recv_service.nim | 14 +++++-- .../waku/node/subscription_manager.nim | 38 ++++++++++--------- logos_delivery/waku/node/waku_node.nim | 2 +- 4 files changed, 37 insertions(+), 24 deletions(-) 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() From 84a42a3eb4f684283c5dc6dd7246b903f2af97ea Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:44:33 +0200 Subject: [PATCH 04/10] refactor(consumers): reuse envelope.hash in archive/filter/REST/FFI; share filter buffer archive.handleMessage and filter.handleMessage gain envelope-native entry points (reusing the precomputed hash) with thin (topic,msg) compat overloads for tests. Filter pushToPeer takes a shared `ref seq[byte]` so the encoded buffer is not deep-copied per subscribed peer. REST messageCacheHandler and the FFI relay handler adopt the envelope; JsonMessageEvent.new reuses the passed-in hash. Co-Authored-By: Claude Fable 5 --- library/events/json_message_event.nim | 14 ++++++-- library/kernel_api/protocols/relay_api.nim | 4 +-- logos_delivery/waku/rest_api/handlers.nim | 4 +-- logos_delivery/waku/waku_archive/archive.nim | 16 ++++++--- .../waku/waku_filter_v2/protocol.nim | 35 ++++++++++++------- 5 files changed, 50 insertions(+), 23 deletions(-) diff --git a/library/events/json_message_event.nim b/library/events/json_message_event.nim index 61278b4fa..ca44300c3 100644 --- a/library/events/json_message_event.nim +++ b/library/events/json_message_event.nim @@ -69,9 +69,15 @@ type JsonMessageEvent* = ref object of JsonEvent messageHash*: string wakuMessage*: JsonMessage -proc new*(T: type JsonMessageEvent, pubSubTopic: string, msg: WakuMessage): T = +proc new*( + T: type JsonMessageEvent, + pubSubTopic: string, + msg: WakuMessage, + msgHash: WakuMessageHash, +): T = # Returns a WakuMessage event as indicated in # https://github.com/vacp2p/rfc/blob/master/content/docs/rfcs/36/README.md#jsonmessageevent-type + # `msgHash` is the precomputed message hash (reused from the inbound envelope). var payload = newSeq[byte](len(msg.payload)) if len(msg.payload) != 0: @@ -85,8 +91,6 @@ proc new*(T: type JsonMessageEvent, pubSubTopic: string, msg: WakuMessage): T = if len(msg.proof) != 0: copyMem(addr proof[0], unsafeAddr msg.proof[0], len(msg.proof)) - let msgHash = computeMessageHash(pubSubTopic, msg) - return JsonMessageEvent( eventType: "message", pubSubTopic: pubSubTopic, @@ -102,5 +106,9 @@ proc new*(T: type JsonMessageEvent, pubSubTopic: string, msg: WakuMessage): T = ), ) +proc new*(T: type JsonMessageEvent, pubSubTopic: string, msg: WakuMessage): T = + ## Convenience overload computing the hash for callers without one. + JsonMessageEvent.new(pubSubTopic, msg, computeMessageHash(pubSubTopic, msg)) + method `$`*(jsonMessage: JsonMessageEvent): string = $(%*jsonMessage) diff --git a/library/kernel_api/protocols/relay_api.nim b/library/kernel_api/protocols/relay_api.nim index 7a1fe446f..a2e23398c 100644 --- a/library/kernel_api/protocols/relay_api.nim +++ b/library/kernel_api/protocols/relay_api.nim @@ -78,9 +78,9 @@ proc waku_relay_subscribe( pubSubTopic: cstring, ) {.ffi.} = proc onReceivedMessage(ctx: ptr FFIContext[LogosDelivery]): WakuRelayHandler = - return proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async.} = + return proc(envelope: WakuEnvelope) {.async.} = callEventCallback(ctx, "onReceivedMessage"): - $JsonMessageEvent.new(pubsubTopic, msg) + $JsonMessageEvent.new(envelope.pubsubTopic, envelope.msg, envelope.hash) ( await ctx.myLib[].waku.relaySubscribe( diff --git a/logos_delivery/waku/rest_api/handlers.nim b/logos_delivery/waku/rest_api/handlers.nim index 520c0518a..df36bdf3e 100644 --- a/logos_delivery/waku/rest_api/handlers.nim +++ b/logos_delivery/waku/rest_api/handlers.nim @@ -34,5 +34,5 @@ proc defaultDiscoveryHandler*( ### Message Cache proc messageCacheHandler*(cache: MessageCache): WakuRelayHandler = - return proc(pubsubTopic: string, msg: WakuMessage): Future[void] {.async, closure.} = - cache.addMessage(pubsubTopic, msg) + return proc(envelope: WakuEnvelope): Future[void] {.async, closure.} = + cache.addMessage(envelope.pubsubTopic, envelope.msg) diff --git a/logos_delivery/waku/waku_archive/archive.nim b/logos_delivery/waku/waku_archive/archive.nim index ea5e2bf02..cac963922 100644 --- a/logos_delivery/waku/waku_archive/archive.nim +++ b/logos_delivery/waku/waku_archive/archive.nim @@ -95,10 +95,10 @@ proc new*( return ok(archive) -proc handleMessage*( - self: WakuArchive, pubsubTopic: PubsubTopic, msg: WakuMessage -) {.async.} = - let msgHash = computeMessageHash(pubsubTopic, msg) +proc handleMessage*(self: WakuArchive, envelope: WakuEnvelope) {.async.} = + let pubsubTopic = envelope.pubsubTopic + let msg = envelope.msg + let msgHash = envelope.hash let msgHashHex = msgHash.to0xHex() trace "handling message", @@ -144,6 +144,14 @@ proc handleMessage*( timestamp = msg.timestamp, insertDuration = insertDuration +proc handleMessage*( + self: WakuArchive, pubsubTopic: PubsubTopic, msg: WakuMessage +) {.async.} = + ## Convenience overload building the envelope (and its hash) for callers that + ## don't already have one (e.g. tests, REST). The relay dispatch path uses the + ## envelope overload directly to avoid re-hashing. + await self.handleMessage(WakuEnvelope.init(pubsubTopic, msg)) + proc syncMessageIngress*( self: WakuArchive, msgHash: WakuMessageHash, diff --git a/logos_delivery/waku/waku_filter_v2/protocol.nim b/logos_delivery/waku/waku_filter_v2/protocol.nim index ab77ce2dc..7bcd90c42 100644 --- a/logos_delivery/waku/waku_filter_v2/protocol.nim +++ b/logos_delivery/waku/waku_filter_v2/protocol.nim @@ -168,8 +168,10 @@ proc handleSubscribeRequest*( return FilterSubscribeResponse.ok(request.requestId) proc pushToPeer( - wf: WakuFilter, peerId: PeerId, buffer: seq[byte] + wf: WakuFilter, peerId: PeerId, buffer: ref seq[byte] ): Future[Result[void, string]] {.async.} = + ## `buffer` is shared by reference across all target peers (encoded once in + ## `pushToPeers`) so no per-peer async-closure copy of the payload happens. info "pushing message to subscribed peer", peerId = shortLog(peerId) let stream = ( @@ -178,21 +180,21 @@ proc pushToPeer( error "pushToPeer failed", error return err("pushToPeer failed: " & $error) - await stream.writeLp(buffer) + await stream.writeLp(buffer[]) info "published successful", peerId = shortLog(peerId), stream waku_service_network_bytes.inc( - amount = buffer.len().int64, labelValues = [WakuFilterPushCodec, "out"] + amount = buffer[].len().int64, labelValues = [WakuFilterPushCodec, "out"] ) return ok() proc pushToPeers( - wf: WakuFilter, peers: seq[PeerId], messagePush: MessagePush + wf: WakuFilter, peers: seq[PeerId], messagePush: MessagePush, msgHash: string ) {.async.} = + ## `msgHash` is the precomputed 0x-hex hash of the pushed message (reused from + ## the inbound envelope; not recomputed here). let targetPeerIds = peers.mapIt(shortLog(it)) - let msgHash = - messagePush.pubsubTopic.computeMessageHash(messagePush.wakuMessage).to0xHex() ## it's also refresh expire of msghash, that's why update cache every time, even if it has a value. if wf.messageCache.put(msgHash, Moment.now()): @@ -210,7 +212,8 @@ proc pushToPeers( target_peer_ids = targetPeerIds, msg_hash = msgHash - let bufferToPublish = messagePush.encode().buffer + let bufferToPublish = new(seq[byte]) + bufferToPublish[] = messagePush.encode().buffer var pushFuts: seq[Future[Result[void, string]]] for peerId in peers: @@ -240,10 +243,10 @@ proc maintainSubscriptions*(wf: WakuFilter) {.async.} = waku_filter_subscriptions.set(wf.subscriptions.peersSubscribed.len.float64) const MessagePushTimeout = 20.seconds -proc handleMessage*( - wf: WakuFilter, pubsubTopic: PubsubTopic, message: WakuMessage -) {.async.} = - let msgHash = computeMessageHash(pubsubTopic, message).to0xHex() +proc handleMessage*(wf: WakuFilter, envelope: WakuEnvelope) {.async.} = + let pubsubTopic = envelope.pubsubTopic + let message = envelope.msg + let msgHash = envelope.hash.to0xHex() info "handling message", pubsubTopic = pubsubTopic, contentTopic = message.contentTopic, msg_hash = msgHash @@ -263,7 +266,7 @@ proc handleMessage*( let messagePush = MessagePush(pubsubTopic: pubsubTopic, wakuMessage: message) - if not await wf.pushToPeers(subscribedPeers, messagePush).withTimeout( + if not await wf.pushToPeers(subscribedPeers, messagePush, msgHash).withTimeout( MessagePushTimeout ): error "timed out pushing message to peers", @@ -287,6 +290,14 @@ proc handleMessage*( # Duration in seconds with millisecond precision floating point waku_filter_handle_message_duration_seconds.observe(handleMessageDurationSec) +proc handleMessage*( + wf: WakuFilter, pubsubTopic: PubsubTopic, message: WakuMessage +) {.async.} = + ## Convenience overload building the envelope (and its hash) for callers that + ## don't already have one (e.g. tests). The relay dispatch path uses the + ## envelope overload directly to avoid re-hashing. + await wf.handleMessage(WakuEnvelope.init(pubsubTopic, message)) + proc initProtocolHandler(wf: WakuFilter) = proc handler(conn: Connection, proto: string) {.async: (raises: [CancelledError]).} = info "filter subscribe request handler triggered", From e1ca47ce33478dabfa18ae0b01c1aa6a52f4f6cf Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:44:40 +0200 Subject: [PATCH 05/10] refactor(apps): adopt WakuEnvelope in relay handlers and message-path bench chat2 and networkmonitor relay handlers take the envelope; the micro/macro benchmark's app handlers follow the new signature (paced/windowed publishing and counters unchanged). Co-Authored-By: Claude Fable 5 --- apps/benchmarks/message_path_bench.nim | 4 ++-- apps/chat2/chat2.nim | 4 +++- apps/networkmonitor/networkmonitor.nim | 6 +++--- 3 files changed, 8 insertions(+), 6 deletions(-) diff --git a/apps/benchmarks/message_path_bench.nim b/apps/benchmarks/message_path_bench.nim index af9333166..e413e44bf 100644 --- a/apps/benchmarks/message_path_bench.nim +++ b/apps/benchmarks/message_path_bench.nim @@ -211,7 +211,7 @@ proc emit(res: ScenarioResult) = # Node setup # --------------------------------------------------------------------------- -proc dummyHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = +proc dummyHandler(envelope: WakuEnvelope) {.async, gcsafe.} = discard proc buildIngestNode( @@ -304,7 +304,7 @@ proc runMacro(shard: PubsubTopic, work: Workload): Future[ScenarioResult] {.asyn times: newSeqOfCap[MonoTime](work.msgs.len), ) - proc countingHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = + proc countingHandler(envelope: WakuEnvelope) {.async, gcsafe.} = arrivals.times.add(getMonoTime()) arrivals.count += 1 if arrivals.count >= arrivals.target: diff --git a/apps/chat2/chat2.nim b/apps/chat2/chat2.nim index 647667b0c..38c0a9740 100644 --- a/apps/chat2/chat2.nim +++ b/apps/chat2/chat2.nim @@ -508,7 +508,9 @@ proc processInput(rfd: AsyncFD, rng: crypto.Rng) {.async.} = # Subscribe to a topic, if relay is mounted if conf.relay: - proc handler(topic: PubsubTopic, msg: WakuMessage): Future[void] {.async, gcsafe.} = + proc handler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic = envelope.pubsubTopic + let msg = envelope.msg trace "Hit subscribe handler", topic if msg.contentTopic == chat.contentTopic: diff --git a/apps/networkmonitor/networkmonitor.nim b/apps/networkmonitor/networkmonitor.nim index a23c778ef..7cc3bbcdc 100644 --- a/apps/networkmonitor/networkmonitor.nim +++ b/apps/networkmonitor/networkmonitor.nim @@ -518,9 +518,9 @@ proc subscribeAndHandleMessages( msgPerContentTopic: ContentTopicMessageTableRef, ) = # handle function - proc handler( - pubsubTopic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc handler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let pubsubTopic = envelope.pubsubTopic + let msg = envelope.msg trace "rx message", pubsubTopic = pubsubTopic, contentTopic = msg.contentTopic # If we reach a table limit size, remove c topics with the least messages. From dab42f85fcc022d70256df5d58741b165b3f34a6 Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:50:35 +0200 Subject: [PATCH 06/10] test: adopt WakuEnvelope in relay handlers across suites Relay app handlers (WakuRelayHandler) now take a WakuEnvelope; handlers that used topic/msg get {.used.} aliases. Archive/filter/store-sync test call sites are unchanged thanks to the compat overloads. Co-Authored-By: Claude Fable 5 --- tests/api/test_api_health.nim | 4 +- tests/api/test_api_receive.nim | 2 +- tests/api/test_api_send.nim | 8 +- tests/api/test_api_subscription.nim | 6 +- tests/node/test_wakunode_health_monitor.nim | 2 +- tests/node/test_wakunode_legacy_lightpush.nim | 6 +- tests/node/test_wakunode_lightpush.nim | 8 +- tests/node/test_wakunode_relay_rln.nim | 12 +-- tests/test_peer_manager.nim | 6 +- tests/test_relay_peer_exchange.nim | 6 +- tests/test_waku_metadata.nim | 2 +- tests/test_wakunode.nim | 6 +- tests/waku_relay/test_protocol.nim | 96 +++++++++---------- tests/waku_relay/test_wakunode_relay.nim | 65 +++++++------ tests/waku_relay/utils.nim | 15 +-- .../test_wakunode_rln_relay.nim | 54 +++++------ tests/waku_rln_relay/utils_offchain.nim | 6 +- tests/wakunode2/test_validators.nim | 6 +- tests/wakunode_rest/test_rest_admin.nim | 6 +- tests/wakunode_rest/test_rest_filter.nim | 18 ++-- tests/wakunode_rest/test_rest_lightpush.nim | 18 ++-- .../test_rest_lightpush_legacy.nim | 18 ++-- tests/wakunode_rest/test_rest_relay.nim | 52 +++++----- 23 files changed, 202 insertions(+), 220 deletions(-) diff --git a/tests/api/test_api_health.nim b/tests/api/test_api_health.nim index 7b493964a..4e298ba5c 100644 --- a/tests/api/test_api_health.nim +++ b/tests/api/test_api_health.nim @@ -21,9 +21,7 @@ const TestTimeout = chronos.seconds(10) const DefaultShard = PubsubTopic("/waku/2/rs/3/0") const TestContentTopic = ContentTopic("/waku/2/default-content/proto") -proc dummyHandler( - topic: PubsubTopic, msg: WakuMessage -): Future[void] {.async, gcsafe.} = +proc dummyHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = discard proc waitForConnectionStatus( diff --git a/tests/api/test_api_receive.nim b/tests/api/test_api_receive.nim index 870f85860..5f291799f 100644 --- a/tests/api/test_api_receive.nim +++ b/tests/api/test_api_receive.nim @@ -107,7 +107,7 @@ proc setupNetwork(testTopic: ContentTopic): Future[TestNetwork] {.async.} = const numShards: uint16 = 1 let shard = PubsubTopic("/waku/2/rs/3/0") - proc dummyHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = + proc dummyHandler(envelope: WakuEnvelope) {.async, gcsafe.} = discard # store node: archive + store + relay, subscribed to the shard diff --git a/tests/api/test_api_send.nim b/tests/api/test_api_send.nim index 8df564b0d..97056b865 100644 --- a/tests/api/test_api_send.nim +++ b/tests/api/test_api_send.nim @@ -209,9 +209,7 @@ suite "Waku API - Send": # Subscribe all relay nodes to the default shard topic const testPubsubTopic = PubsubTopic("/waku/2/rs/3/0") - proc dummyHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc dummyHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = discard relayNode1.subscribe((kind: PubsubSub, topic: testPubsubTopic), dummyHandler).isOkOr: @@ -476,9 +474,7 @@ suite "Waku API - Send": await fakeLightpushNode.mountLibp2pPing() await fakeLightpushNode.start() let fakeLightpushNodePeerInfo = fakeLightpushNode.peerInfo.toRemotePeerInfo() - proc dummyHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc dummyHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = discard fakeLightpushNode.subscribe( diff --git a/tests/api/test_api_subscription.nim b/tests/api/test_api_subscription.nim index 90d160bc3..7f777a30a 100644 --- a/tests/api/test_api_subscription.nim +++ b/tests/api/test_api_subscription.nim @@ -108,7 +108,7 @@ proc setupNetwork( net.publisherPeerInfo = net.publisher.peerInfo.toRemotePeerInfo() - proc dummyHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = + proc dummyHandler(envelope: WakuEnvelope) {.async, gcsafe.} = discard var shards: seq[PubsubTopic] @@ -617,7 +617,7 @@ suite "Messaging API, SubscriptionManager": let numShards: uint16 = 1 let shards = @[PubsubTopic("/waku/2/rs/3/0")] - proc dummyHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = + proc dummyHandler(envelope: WakuEnvelope) {.async, gcsafe.} = discard var publisher: WakuNode @@ -726,7 +726,7 @@ suite "Messaging API, SubscriptionManager": let numShards: uint16 = 1 let shards = @[PubsubTopic("/waku/2/rs/3/0")] - proc dummyHandler(topic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = + proc dummyHandler(envelope: WakuEnvelope) {.async, gcsafe.} = discard var publisher: WakuNode diff --git a/tests/node/test_wakunode_health_monitor.nim b/tests/node/test_wakunode_health_monitor.nim index af90a9be0..c2f69a19e 100644 --- a/tests/node/test_wakunode_health_monitor.nim +++ b/tests/node/test_wakunode_health_monitor.nim @@ -171,7 +171,7 @@ suite "Health Monitor - events": await nodeA.connectToNodes(@[nodeB.switch.peerInfo.toRemotePeerInfo()]) - proc dummyHandler(topic: PubsubTopic, msg: WakuMessage): Future[void] {.async.} = + proc dummyHandler(envelope: WakuEnvelope): Future[void] {.async.} = discard nodeA.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), dummyHandler).expect( diff --git a/tests/node/test_wakunode_legacy_lightpush.nim b/tests/node/test_wakunode_legacy_lightpush.nim index aec37e18c..33f4f612b 100644 --- a/tests/node/test_wakunode_legacy_lightpush.nim +++ b/tests/node/test_wakunode_legacy_lightpush.nim @@ -303,9 +303,9 @@ suite "Waku Legacy Lightpush message delivery": const CustomPubsubTopic = "/waku/2/rs/0/1" let message = fakeWakuMessage() var completionFutRelay = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == CustomPubsubTopic msg == message diff --git a/tests/node/test_wakunode_lightpush.nim b/tests/node/test_wakunode_lightpush.nim index f13cbcaab..f2605c592 100644 --- a/tests/node/test_wakunode_lightpush.nim +++ b/tests/node/test_wakunode_lightpush.nim @@ -387,12 +387,10 @@ suite "Waku Lightpush message delivery": let message = fakeWakuMessage() var completionFutRelay = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = check: - topic == CustomPubsubTopic - msg == message + envelope.pubsubTopic == CustomPubsubTopic + envelope.msg == message completionFutRelay.complete(true) destNode.subscribe((kind: PubsubSub, topic: CustomPubsubTopic), relayHandler).isOkOr: diff --git a/tests/node/test_wakunode_relay_rln.nim b/tests/node/test_wakunode_relay_rln.nim index 76074bb6d..a0dcc4254 100644 --- a/tests/node/test_wakunode_relay_rln.nim +++ b/tests/node/test_wakunode_relay_rln.nim @@ -232,9 +232,9 @@ suite "Waku RlnRelay - End to End - Static": # Register Relay Handler var completionFut = newPushHandlerFuture() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg if topic == pubsubTopic: completionFut.complete((topic, msg)) @@ -325,9 +325,9 @@ suite "Waku RlnRelay - End to End - Static": # Register Relay Handler var completionFut = newPushHandlerFuture() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg if topic == pubsubTopic: completionFut.complete((topic, msg)) diff --git a/tests/test_peer_manager.nim b/tests/test_peer_manager.nim index b364fc8c3..ccec92a1b 100644 --- a/tests/test_peer_manager.nim +++ b/tests/test_peer_manager.nim @@ -668,9 +668,9 @@ procSuite "Peer Manager": await allFutures(nodes.mapIt(it.mountRelay())) await allFutures(nodes.mapIt(it.start())) - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.millis) let topic = "/waku/2/rs/0/0" diff --git a/tests/test_relay_peer_exchange.nim b/tests/test_relay_peer_exchange.nim index 058576d4c..79e48b340 100644 --- a/tests/test_relay_peer_exchange.nim +++ b/tests/test_relay_peer_exchange.nim @@ -91,9 +91,9 @@ procSuite "Relay (GossipSub) Peer Exchange": await allFutures([node1.start(), node2.start(), node3.start()]) # The three nodes should be subscribed to the same shard - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node1.subscribe((kind: PubsubSub, topic: $DefaultRelayShard), simpleHandler).isOkOr: diff --git a/tests/test_waku_metadata.nim b/tests/test_waku_metadata.nim index 1ae4ce7b7..0f06e327e 100644 --- a/tests/test_waku_metadata.nim +++ b/tests/test_waku_metadata.nim @@ -45,7 +45,7 @@ procSuite "Waku Metadata Protocol": # Subscribe to topics on node1 - relay will track these and metadata will report them let noOpHandler: WakuRelayHandler = proc( - pubsubTopic: PubsubTopic, message: WakuMessage + envelope: WakuEnvelope ): Future[void] {.async.} = discard diff --git a/tests/test_wakunode.nim b/tests/test_wakunode.nim index 4279d9066..645d78466 100644 --- a/tests/test_wakunode.nim +++ b/tests/test_wakunode.nim @@ -58,9 +58,9 @@ suite "WakuNode": await node1.connectToNodes(@[node2.switch.peerInfo.toRemotePeerInfo()]) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard msg.contentTopic == contentTopic diff --git a/tests/waku_relay/test_protocol.nim b/tests/waku_relay/test_protocol.nim index a2042627e..a8485d205 100644 --- a/tests/waku_relay/test_protocol.nim +++ b/tests/waku_relay/test_protocol.nim @@ -60,10 +60,10 @@ suite "Waku Relay": messageSeq = @[] handlerFuture = newPushHandlerFuture() simpleFutureHandler = proc( - topic: PubsubTopic, msg: WakuMessage + envelope: WakuEnvelope ): Future[void] {.async, closure, gcsafe.} = - messageSeq.add((topic, msg)) - handlerFuture.complete((topic, msg)) + messageSeq.add((envelope.pubsubTopic, envelope.msg)) + handlerFuture.complete((envelope.pubsubTopic, envelope.msg)) switch = newTestSwitch() peerManager = PeerManager.new(switch) @@ -123,9 +123,9 @@ suite "Waku Relay": check await peerManager.connectPeer(otherRemotePeerInfo) var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) # When subscribing the second node to the Pubsub Topic @@ -184,9 +184,9 @@ suite "Waku Relay": check await peerManager.connectPeer(otherRemotePeerInfo) var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) # When subscribing both nodes to the same Pubsub Topic @@ -257,9 +257,9 @@ suite "Waku Relay": # Given the subscription is refreshed var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) node.subscribe(pubsubTopic, otherSimpleFutureHandler) @@ -303,9 +303,9 @@ suite "Waku Relay": check await peerManager.connectPeer(otherRemotePeerInfo) var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) otherNode.addValidator(len4Validator) @@ -393,9 +393,9 @@ suite "Waku Relay": check await peerManager.connectPeer(otherRemotePeerInfo) var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) node.subscribe(pubsubTopic, simpleFutureHandler) @@ -477,37 +477,35 @@ suite "Waku Relay": # Given the first node is subscribed to two pubsub topics var handlerFuture2 = newPushHandlerFuture() - proc simpleFutureHandler2( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = - handlerFuture2.complete((topic, message)) + proc simpleFutureHandler2(envelope: WakuEnvelope) {.async, gcsafe.} = + handlerFuture2.complete((envelope.pubsubTopic, envelope.msg)) node.subscribe(pubsubTopic, simpleFutureHandler) node.subscribe(pubsubTopicB, simpleFutureHandler2) # Given the other nodes are subscribed to two pubsub topics var otherHandlerFuture1 = newPushHandlerFuture() - proc otherSimpleFutureHandler1( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler1(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture1.complete((topic, message)) var otherHandlerFuture2 = newPushHandlerFuture() - proc otherSimpleFutureHandler2( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler2(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture2.complete((topic, message)) var anotherHandlerFuture1 = newPushHandlerFuture() - proc anotherSimpleFutureHandler1( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc anotherSimpleFutureHandler1(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg anotherHandlerFuture1.complete((topic, message)) var anotherHandlerFuture2 = newPushHandlerFuture() - proc anotherSimpleFutureHandler2( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc anotherSimpleFutureHandler2(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg anotherHandlerFuture2.complete((topic, message)) otherNode.subscribe(pubsubTopic, otherSimpleFutureHandler1) @@ -870,9 +868,9 @@ suite "Waku Relay": # Given both are subscribed to the same pubsub topic var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) otherNode.subscribe(pubsubTopic, otherSimpleFutureHandler) @@ -1036,9 +1034,9 @@ suite "Waku Relay": # Given both are subscribed to the same pubsub topic var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) otherNode.subscribe(pubsubTopic, otherSimpleFutureHandler) @@ -1169,17 +1167,17 @@ suite "Waku Relay": # Create a different handler than the default to include messages in a seq var thisHandlerFuture = newPushHandlerFuture() var thisMessageSeq: seq[(PubsubTopic, WakuMessage)] = @[] - proc thisSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc thisSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg thisMessageSeq.add((topic, message)) thisHandlerFuture.complete((topic, message)) var otherHandlerFuture = newPushHandlerFuture() var otherMessageSeq: seq[(PubsubTopic, WakuMessage)] = @[] - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherMessageSeq.add((topic, message)) otherHandlerFuture.complete((topic, message)) @@ -1252,9 +1250,9 @@ suite "Waku Relay": # Given both are subscribed to the same pubsub topic var otherHandlerFuture = newPushHandlerFuture() - proc otherSimpleFutureHandler( - topic: PubsubTopic, message: WakuMessage - ) {.async, gcsafe.} = + proc otherSimpleFutureHandler(envelope: WakuEnvelope) {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let message {.used.} = envelope.msg otherHandlerFuture.complete((topic, message)) otherNode.subscribe(pubsubTopic, otherSimpleFutureHandler) diff --git a/tests/waku_relay/test_wakunode_relay.nim b/tests/waku_relay/test_wakunode_relay.nim index 6ec622f2a..7ad27bd7b 100644 --- a/tests/waku_relay/test_wakunode_relay.nim +++ b/tests/waku_relay/test_wakunode_relay.nim @@ -87,9 +87,9 @@ suite "WakuNode - Relay": ) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard msg.contentTopic == contentTopic @@ -97,9 +97,9 @@ suite "WakuNode - Relay": msg.timestamp > 0 completionFut.complete(true) - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) ## node1 and node2 explicitly subscribe to the same shard as node3 @@ -189,9 +189,9 @@ suite "WakuNode - Relay": node2.wakuRelay.addValidator(validator) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard # check that only messages with contentTopic1 is relayed (but not contentTopic2) @@ -199,9 +199,9 @@ suite "WakuNode - Relay": # relay handler is called completionFut.complete(true) - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) ## node1 and node2 explicitly subscribe to the same shard as node3 @@ -295,9 +295,9 @@ suite "WakuNode - Relay": await node1.connectToNodes(@[node2.switch.peerInfo.toRemotePeerInfo()]) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard msg.contentTopic == contentTopic @@ -347,9 +347,9 @@ suite "WakuNode - Relay": await node1.connectToNodes(@[node2.switch.peerInfo.toRemotePeerInfo()]) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard msg.contentTopic == contentTopic @@ -408,9 +408,9 @@ suite "WakuNode - Relay": await node1.connectToNodes(@[node2.switch.peerInfo.toRemotePeerInfo()]) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard msg.contentTopic == contentTopic @@ -467,9 +467,9 @@ suite "WakuNode - Relay": await node1.connectToNodes(@[node2.switch.peerInfo.toRemotePeerInfo()]) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard msg.contentTopic == contentTopic @@ -527,9 +527,9 @@ suite "WakuNode - Relay": await node1.connectToNodes(@[node2.switch.peerInfo.toRemotePeerInfo()]) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg check: topic == $shard msg.contentTopic == contentTopic @@ -563,9 +563,9 @@ suite "WakuNode - Relay": await allFutures(nodes.mapIt(it.start())) await allFutures(nodes.mapIt(it.mountRelay())) - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) # subscribe all nodes to a topic @@ -634,10 +634,9 @@ suite "WakuNode - Relay": contentTopicB = ContentTopic("/waku/2/default-content1/proto") contentTopicC = ContentTopic("/waku/2/default-content2/proto") handler: WakuRelayHandler = proc( - pubsubTopic: PubsubTopic, message: WakuMessage + envelope: WakuEnvelope ): Future[void] {.gcsafe, raises: [Defect].} = - discard pubsubTopic - discard message + discard envelope assert shard == node.wakuAutoSharding.get().getShard(contentTopicA).expect("Valid Topic"), "topic must use the same shard" diff --git a/tests/waku_relay/utils.nim b/tests/waku_relay/utils.nim index 663d06c18..85c8a8fda 100644 --- a/tests/waku_relay/utils.nim +++ b/tests/waku_relay/utils.nim @@ -21,7 +21,7 @@ import proc noopRawHandler*(): WakuRelayHandler = var handler: WakuRelayHandler - handler = proc(topic: PubsubTopic, msg: WakuMessage): Future[void] {.async, gcsafe.} = + handler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = discard handler @@ -39,11 +39,8 @@ proc subscribeToContentTopicWithHandler*( node: WakuNode, contentTopic: string ): Future[bool] = var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = - if topic == topic: - completionFut.complete(true) + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + completionFut.complete(true) (node.subscribe((kind: ContentSub, topic: contentTopic), relayHandler)).isOkOr: error "Failed to subscribe to content topic", error @@ -52,10 +49,8 @@ proc subscribeToContentTopicWithHandler*( proc subscribeCompletionHandler*(node: WakuNode, pubsubTopic: string): Future[bool] = var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = - if topic == pubsubTopic: + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + if envelope.pubsubTopic == pubsubTopic: completionFut.complete(true) (node.subscribe((kind: PubsubSub, topic: pubsubTopic), relayHandler)).isOkOr: diff --git a/tests/waku_rln_relay/test_wakunode_rln_relay.nim b/tests/waku_rln_relay/test_wakunode_rln_relay.nim index 4c02b4cbd..9268f8b87 100644 --- a/tests/waku_rln_relay/test_wakunode_rln_relay.nim +++ b/tests/waku_rln_relay/test_wakunode_rln_relay.nim @@ -107,16 +107,16 @@ procSuite "WakuNode - RLN relay": await node3.connectToNodes(@[node2.switch.peerInfo.toRemotePeerInfo()]) var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg info "The received topic:", topic if topic == DefaultPubsubTopic: completionFut.complete(true) - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node1.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -224,18 +224,18 @@ procSuite "WakuNode - RLN relay": var rxMessagesTopic1 = 0 var rxMessagesTopic2 = 0 - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg info "relayHandler. The received topic:", topic if topic == $shards[0]: rxMessagesTopic1 = rxMessagesTopic1 + 1 elif topic == $shards[1]: rxMessagesTopic2 = rxMessagesTopic2 + 1 - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node1.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -366,16 +366,16 @@ procSuite "WakuNode - RLN relay": # define a custom relay handler var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg info "The received topic:", topic if topic == DefaultPubsubTopic: completionFut.complete(true) - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node1.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -526,9 +526,9 @@ procSuite "WakuNode - RLN relay": var completionFut2 = newFuture[bool]() var completionFut3 = newFuture[bool]() var completionFut4 = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg info "The received topic:", topic if topic == DefaultPubsubTopic: if msg == wm1: @@ -540,9 +540,9 @@ procSuite "WakuNode - RLN relay": if msg.payload == wm4.payload: completionFut4.complete(true) - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node1.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -649,9 +649,9 @@ procSuite "WakuNode - RLN relay": completionFut4 = newFuture[bool]() completionFut5 = newFuture[bool]() completionFut6 = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg info "The received topic:", topic if topic == DefaultPubsubTopic: if msg == wm1: diff --git a/tests/waku_rln_relay/utils_offchain.nim b/tests/waku_rln_relay/utils_offchain.nim index 9976ec783..478fbed9f 100644 --- a/tests/waku_rln_relay/utils_offchain.nim +++ b/tests/waku_rln_relay/utils_offchain.nim @@ -34,9 +34,9 @@ proc setupRelayWithStaticRln*( proc subscribeCompletionHandler*(node: WakuNode, pubsubTopic: string): Future[bool] = var completionFut = newFuture[bool]() - proc relayHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc relayHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg if topic == pubsubTopic: completionFut.complete(true) diff --git a/tests/wakunode2/test_validators.nim b/tests/wakunode2/test_validators.nim index fc4c2cf57..821ca6e54 100644 --- a/tests/wakunode2/test_validators.nim +++ b/tests/wakunode2/test_validators.nim @@ -61,7 +61,7 @@ suite "WakuNode2 - Validators": await sleepAsync(500.millis) var msgReceived = 0 - proc handler(pubsubTopic: PubsubTopic, data: WakuMessage) {.async, gcsafe.} = + proc handler(envelope: WakuEnvelope) {.async, gcsafe.} = msgReceived += 1 # Subscribe all nodes to the same topic/handler @@ -148,7 +148,7 @@ suite "WakuNode2 - Validators": require connOk var msgReceived = 0 - proc handler(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = + proc handler(envelope: WakuEnvelope) {.async, gcsafe.} = msgReceived += 1 # Connection triggers different actions, wait for them @@ -279,7 +279,7 @@ suite "WakuNode2 - Validators": await allFutures(nodes.mapIt(it.mountRelay())) var msgReceived = 0 - proc handler(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async, gcsafe.} = + proc handler(envelope: WakuEnvelope) {.async, gcsafe.} = msgReceived += 1 # Subscribe all nodes to the same topic/handler diff --git a/tests/wakunode_rest/test_rest_admin.nim b/tests/wakunode_rest/test_rest_admin.nim index 81ef7e6ea..daf6771c6 100644 --- a/tests/wakunode_rest/test_rest_admin.nim +++ b/tests/wakunode_rest/test_rest_admin.nim @@ -60,9 +60,9 @@ suite "Waku v2 Rest API - Admin": ) # The three nodes should be subscribed to the same shard - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) let shard = RelayShard(clusterId: clusterId, shardId: 5) diff --git a/tests/wakunode_rest/test_rest_filter.nim b/tests/wakunode_rest/test_rest_filter.nim index b20e67e59..b762d8e3e 100644 --- a/tests/wakunode_rest/test_rest_filter.nim +++ b/tests/wakunode_rest/test_rest_filter.nim @@ -278,9 +278,9 @@ suite "Waku v2 Rest API - Filter V2": restFilterTest = await RestFilterTest.init() subPeerId = restFilterTest.subscriberNode.peerInfo.toRemotePeerInfo().peerId - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restFilterTest.messageCache.pubsubSubscribe(DefaultPubsubTopic) @@ -333,9 +333,9 @@ suite "Waku v2 Rest API - Filter V2": # setup filter service and client node let restFilterTest = await RestFilterTest.init() let subPeerId = restFilterTest.subscriberNode.peerInfo.toRemotePeerInfo().peerId - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restFilterTest.serviceNode.subscribe( @@ -412,9 +412,9 @@ suite "Waku v2 Rest API - Filter V2": # setup filter service and client node let restFilterTest = await RestFilterTest.init() let subPeerId = restFilterTest.subscriberNode.peerInfo.toRemotePeerInfo().peerId - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restFilterTest.serviceNode.subscribe( diff --git a/tests/wakunode_rest/test_rest_lightpush.nim b/tests/wakunode_rest/test_rest_lightpush.nim index ff602328a..915096612 100644 --- a/tests/wakunode_rest/test_rest_lightpush.nim +++ b/tests/wakunode_rest/test_rest_lightpush.nim @@ -129,9 +129,9 @@ suite "Waku v2 Rest API - lightpush": # Given let restLightPushTest = await RestLightPushTest.init() - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restLightPushTest.consumerNode.subscribe( @@ -168,9 +168,9 @@ suite "Waku v2 Rest API - lightpush": asyncTest "Push message bad-request": # Given let restLightPushTest = await RestLightPushTest.init() - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restLightPushTest.serviceNode.subscribe( @@ -230,9 +230,9 @@ suite "Waku v2 Rest API - lightpush": let budgetCap = 3 let tokenPeriod = 500.millis let restLightPushTest = await RestLightPushTest.init((budgetCap, tokenPeriod)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restLightPushTest.consumerNode.subscribe( diff --git a/tests/wakunode_rest/test_rest_lightpush_legacy.nim b/tests/wakunode_rest/test_rest_lightpush_legacy.nim index 5f29146c8..5df6cc8b2 100644 --- a/tests/wakunode_rest/test_rest_lightpush_legacy.nim +++ b/tests/wakunode_rest/test_rest_lightpush_legacy.nim @@ -123,9 +123,9 @@ suite "Waku v2 Rest API - lightpush": asyncTest "Push message request": # Given let restLightPushTest = await RestLightPushTest.init() - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restLightPushTest.consumerNode.subscribe( @@ -162,9 +162,9 @@ suite "Waku v2 Rest API - lightpush": asyncTest "Push message bad-request": # Given let restLightPushTest = await RestLightPushTest.init() - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restLightPushTest.serviceNode.subscribe( @@ -227,9 +227,9 @@ suite "Waku v2 Rest API - lightpush": let budgetCap = 3 let tokenPeriod = 500.millis let restLightPushTest = await RestLightPushTest.init((budgetCap, tokenPeriod)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) restLightPushTest.consumerNode.subscribe( diff --git a/tests/wakunode_rest/test_rest_relay.nim b/tests/wakunode_rest/test_rest_relay.nim index b59dc463d..6ae9850c1 100644 --- a/tests/wakunode_rest/test_rest_relay.nim +++ b/tests/wakunode_rest/test_rest_relay.nim @@ -126,9 +126,9 @@ suite "Waku v2 Rest API - Relay": (await node.mountRelay()).isOkOr: assert false, "Failed to mount relay" - proc simpleHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc simpleHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) for shard in @[$shard0, $shard1, $shard2, $shard3, $shard4]: @@ -292,9 +292,9 @@ suite "Waku v2 Rest API - Relay": let client = newRestHttpClient(initTAddress(restAddress, restPort)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -512,9 +512,7 @@ suite "Waku v2 Rest API - Relay": await meshNode.setRlnValidator(wakuRlnConfig) await meshNode.start() const testPubsubTopic = PubsubTopic("/waku/2/rs/1/0") - proc dummyHandler( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + proc dummyHandler(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = discard meshNode.subscribe((kind: ContentSub, topic: DefaultContentTopic), dummyHandler).isOkOr: @@ -562,9 +560,9 @@ suite "Waku v2 Rest API - Relay": let client = newRestHttpClient(initTAddress(restAddress, restPort)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node.subscribe((kind: ContentSub, topic: DefaultContentTopic), simpleHandler).isOkOr: @@ -691,9 +689,9 @@ suite "Waku v2 Rest API - Relay": let client = newRestHttpClient(initTAddress(restAddress, restPort)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -763,9 +761,9 @@ suite "Waku v2 Rest API - Relay": let client = newRestHttpClient(initTAddress(restAddress, restPort)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -839,9 +837,9 @@ suite "Waku v2 Rest API - Relay": restServer.start() let client = newRestHttpClient(initTAddress(restAddress, restPort)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node.subscribe((kind: PubsubSub, topic: DefaultPubsubTopic), simpleHandler).isOkOr: @@ -897,9 +895,9 @@ suite "Waku v2 Rest API - Relay": 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.} = + let meshHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg discard meshNode.subscribe((kind: ContentSub, topic: DefaultContentTopic), meshHandler).isOkOr: assert false, "Failed to subscribe mesh node" @@ -946,9 +944,9 @@ suite "Waku v2 Rest API - Relay": restServer.start() let client = newRestHttpClient(initTAddress(restAddress, restPort)) - let simpleHandler = proc( - topic: PubsubTopic, msg: WakuMessage - ): Future[void] {.async, gcsafe.} = + let simpleHandler = proc(envelope: WakuEnvelope): Future[void] {.async, gcsafe.} = + let topic {.used.} = envelope.pubsubTopic + let msg {.used.} = envelope.msg await sleepAsync(0.milliseconds) node.subscribe((kind: ContentSub, topic: DefaultContentTopic), simpleHandler).isOkOr: From a6f1065a05dee4ed537f977702a4033e89371b4e Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:52:52 +0200 Subject: [PATCH 07/10] =?UTF-8?q?docs(bench):=20Phase=203=20results=20?= =?UTF-8?q?=E2=80=94=20micro=201.0=20hash/msg=20@=20~2220=20msg/s,=20macro?= =?UTF-8?q?=203/3=20@=20356=20msg/s?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Receiver-side gate met: micro isolates it at 2.0 decodes / 1.0 hash per message; macro aggregate 3/3 decomposes to node B = 2 decodes / 1 hash. Co-Authored-By: Claude Fable 5 --- docs/analysis/bench_baseline.md | 67 +++++++++++++++++++++++++++++++++ 1 file changed, 67 insertions(+) diff --git a/docs/analysis/bench_baseline.md b/docs/analysis/bench_baseline.md index 92e79bc88..fa19759d1 100644 --- a/docs/analysis/bench_baseline.md +++ b/docs/analysis/bench_baseline.md @@ -164,3 +164,70 @@ force-copy path exists. | 2 | Mutation-site classification produced | **PASS** — `docs/analysis/phase2_mutation_audit.md` | | 3 | Alloc volume down ≥ 40 %, decode/hash unchanged | **PARTIAL** — decode/hash byte-exact unchanged (**PASS**); alloc-volume drop **not demonstrable** via this harness's retained-memory metric under refc (technical reason above); no residual deep-copy bug found | | 4 | `detect_changes` reviewed; committed (no push) | **PASS** (gitnexus skipped per environment; grep-based mutation sweep committed instead) | + +--- + +# Phase 3 — `WakuEnvelope`: decode once, hash once + +Re-run of the same harness after Phase 3 (`WakuEnvelope` threaded through relay +dispatch; redundant decodes/hashes removed; filter push buffer shared). Same +machine/toolchain/workload (seed 42, 10/50/150 kB @ 25/50/25 %, N=1000 + 100 +warmup), same build defines (`-d:msgPathCounters -d:chronicles_log_level=ERROR`, +`--passL:librln_v2.0.2.a --passL:-lm`), `--mm:refc`. + +Benchmarked at commit `dab42f85`. Observer decode/hash sites are gated behind +`enabledLogLevel <= LogLevel.DEBUG`; the build is at `ERROR`, so they are +compiled out. INFO would give byte-identical counters (the gate threshold is +DEBUG). + +## Results (CSV) + +``` +scenario,msgs,payload_profile,wall_ms,msg_per_s,ns_per_msg_p50,ns_per_msg_p99,decodes_per_msg,hashes_per_msg,hashed_MB,decoded_MB,occupied_mem_delta_MB,gc_collections +micro_run1,1000,10/50/150kB@25/50/25,449.891,2222.8,346625,1040333,2.000,1.000,65.00,130.11,0.04,7 +micro_run2,1000,10/50/150kB@25/50/25,451.472,2215.0,336166,1287166,2.000,1.000,65.00,130.11,1.19,8 +macro,1000,10/50/150kB@25/50/25,2809.200,356.0,2285375,7529458,3.000,3.000,195.00,195.16,137.54,10 +``` +Micro determinism: run1=2222.8 msg/s, run2=2215.0 msg/s → **variance 0.35 %**. + +## Phase 1 → Phase 2 → Phase 3 + +| Scenario | Metric | Phase 1 | Phase 2 | Phase 3 | +|---|---|---|---|---| +| micro | msg/s | 1206.9 / 1257.6 | 1235.7 / 1252.4 | **2222.8 / 2215.0** (~+78 %) | +| micro | decodes/msg | 2.000 | 2.000 | **2.000** | +| micro | hashes/msg | 2.000 | 2.000 | **1.000** | +| micro | hashed_MB | 130.00 | 130.00 | **65.00** | +| macro | msg/s | 229.0 | 230.2 | **356.0** (~+55 %) | +| macro | decodes/msg | 6.000 | 6.000 | **3.000** | +| macro | hashes/msg | 5.000 | 5.000 | **3.000** | +| macro | hashed_MB / decoded_MB | 325.00 / 390.33 | 325.00 / 390.33 | **195.00 / 195.16** | + +## Acceptance gate — receiver-side decodes ≤ 2, hashes == 1 + +The counter is process-global (publisher node A + receiver node B share it). + +**Micro directly isolates the receiver path** (single node: ordered validator + +topicHandler + full dispatch chain trace/filter/archive/sync/internal, no +publisher leg): **2.000 decodes, 1.000 hashes per message** → gate met +(decodes ≤ 2 ✓, hashes == 1 ✓). The one hash is the envelope construction in the +relay topic handler; archive and store-sync reuse `envelope.hash`. + +**Macro aggregate 3.0 decodes / 3.0 hashes** decomposes as: + +| Counter | Node A (publisher, relay-only, self-subscribed, `triggerSelf`) | Node B (receiver) | Total | +|---|---|---|---| +| decode | 1 (own topicHandler via triggerSelf) | 2 (ordered validator + topicHandler) | **3** | +| hash | 2 (`publish` leg + own envelope) | 1 (envelope; archive & store-sync reuse it) | **3** | + +Receiver side (node B) = **2 decodes, 1 hash** → gate met. The publisher's +contribution is the encode-side `publish` hash plus its own self-delivered +envelope (triggerSelf), not receiver overhead. Aggregate 6.0/5.0 → 3.0/3.0 +lands within the expected "~≤3.0 decodes / ~2.0 hashes" band (hashes at 3.0 +because node A both publishes *and* self-receives; the pure receiver leg is 1.0). + +`computeMessageHash` recomputations removed on the inbound relay path: +`waku_relay/protocol.nim` onValidated/onSend (gated), `waku_archive/archive.nim` +(`envelope.hash`), `waku_store_sync/reconciliation.nim` (hash overload), +`waku_filter_v2/protocol.nim` `pushToPeers`, plus the FFI JSON event and +`recv_service` (reused via `MessageSeenEvent`). From be6bfed6192d05ce7c6f4b3ed0621e67229174fc Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 09:16:11 +0200 Subject: [PATCH 08/10] docs(bench): same-session pre-phase-2 vs post-phase-3 speed+GC comparison Back-to-back A/B on one machine: micro +85.9%, macro +56.6% msg/s; decode/hash byte volumes halved. Peak heap (getMaxMem) -21%/-20.7%. GC-sampling series shows retained-memory slope is workload-driven and unchanged between versions; the transient-copy reduction is not observable via occupied/total sawtooth under refc (deterministic frees, reused arena) and manifests as a lower absolute heap level, not a gentler slope. Co-Authored-By: Claude Fable 5 --- docs/analysis/speed_gain_pre2_post3.md | 274 +++++++++++++++++++++++++ 1 file changed, 274 insertions(+) create mode 100644 docs/analysis/speed_gain_pre2_post3.md diff --git a/docs/analysis/speed_gain_pre2_post3.md b/docs/analysis/speed_gain_pre2_post3.md new file mode 100644 index 000000000..35b2af866 --- /dev/null +++ b/docs/analysis/speed_gain_pre2_post3.md @@ -0,0 +1,274 @@ +# Speed gain: pre-phase-2 vs post-phase-3 (same-session A/B) + +Rigorous before/after of the no-copy refactor chain, both branches built and run +**back-to-back on the same machine, in the same session**, with identical build +flags and an identical (temporary, uncommitted) GC-sampling patch to +`apps/benchmarks/message_path_bench.nim`. + +## Methodology + +| Item | Value | +|---|---| +| Branch A ("pre-phase-2") | `experimental/chore-nocopy-e2e-perf-harness` @ `2a4da679` | +| Branch B ("post-phase-3") | `experimental/chore-nocopy-wakuenvelop` @ `a6f1065a` | +| Machine / OS | Apple M4, macOS (arm64) | +| Nim | 2.2.4, `--mm:refc` | +| Build flags (both) | `--mm:refc --cpu:arm64 --passC:"-arch arm64" --passL:"-arch arm64" -d:msgPathCounters -d:chronicles_log_level=ERROR --passL:librln_v2.0.2.a --passL:-lm` (copied from the `benchMessagePath` nimble task) | +| Workload (both) | seed 42, payload mix 10/50/150 kB @ 25/50/25 %, N=1000 measured + 100 warmup | +| Run order | A built+run (2 passes) → B built+run (3 passes; B pass 1 discarded, micro variance 5.51 % under a load spike to 13.8) | +| Machine load | A passes: load avg ~3.3. B passes: elevated (8.9–13.8) — noted; the two accepted B passes have micro variance 3.34 % / 4.15 % (< 5 %). | + +**Sampling design (temporary patch, identical on both branches — only the relay +handler signature differs between branches, which the sampling code does not +touch):** +- Micro loop and macro receive loop: every 50 messages record + `(msg_index, getOccupiedMem(), getTotalMem())` into a pre-allocated buffer + (no per-sample heap churn). +- End of each scenario: `getMaxMem()` (process peak heap — the key single + number) and `GC_getStatistics()`. +- Existing `gc_collections` counter retained. +- All raw outputs saved outside the repo; the patch was `git checkout --` + discarded before switching branches (verified clean each time). + +The instrumented builds are byte-identical in decode/hash counters to the +committed `bench_baseline.md` numbers, confirming the patch did not perturb the +measured path. + +## Throughput — pre vs post + +Micro = single-node isolated receiver path; macro = two loopback nodes, +publisher + archiving receiver. Micro msg/s is the mean of all accepted +micro runs (A: 4 values across 2 passes; B: 4 values across 2 passes). + +| Scenario | Metric | Pre (A) | Post (B) | Gain | +|---|---|---|---|---| +| micro | msg/s (mean) | 1222.8 | 2273.7 | **+85.9 %** | +| micro | ns/msg p50 | ~630 k | ~338 k | −46 % | +| micro | ns/msg p99 | ~1.94 M | ~1.09 M | −44 % | +| micro | decodes/msg | 2.000 | 2.000 | unchanged | +| micro | hashes/msg | 2.000 | **1.000** | −50 % | +| micro | hashed_MB | 130.00 | **65.00** | −50 % | +| macro | msg/s (mean) | 225.45 | 353.1 | **+56.6 %** | +| macro | ns/msg p50 | ~3.55 M | ~2.29 M | −36 % | +| macro | ns/msg p99 | ~10.7 M | ~7.5 M | −30 % | +| macro | decodes/msg | 6.000 | **3.000** | −50 % | +| macro | hashes/msg | 5.000 | **3.000** | −40 % | +| macro | decoded_MB / hashed_MB | 390.33 / 325.00 | **195.16 / 195.00** | ~−50 % | + +Raw CSV rows (accepted passes): + +``` +# A (pre-phase-2) — pass 1 / pass 2 +micro_run1,1000,834.793ms,1197.9,p50=640000,p99=1980291,dec=2.000,hash=2.000 +micro_run2,1000,803.157ms,1245.1,p50=618708,p99=1928250,dec=2.000,hash=2.000 +micro_run1,1000,835.163ms,1197.4,p50=642250,p99=1939583,dec=2.000,hash=2.000 +micro_run2,1000,799.486ms,1250.8,p50=621750,p99=1866667,dec=2.000,hash=2.000 +macro,1000,4431.571ms,225.7,p50=3552333,p99=10703625,dec=6.000,hash=5.000 +macro,1000,4441.265ms,225.2,p50=3580583,p99=10964584,dec=6.000,hash=5.000 +# B (post-phase-3) — pass 2 / pass 3 (pass 1 discarded, variance>5% under load) +micro_run1,1000,451.768ms,2213.5,p50=345792,p99=1138042,dec=2.000,hash=1.000 +micro_run2,1000,436.666ms,2290.1,p50=335709,p99=1045417,dec=2.000,hash=1.000 +micro_run1,1000,445.048ms,2246.9,p50=340917,p99=1092750,dec=2.000,hash=1.000 +micro_run2,1000,426.569ms,2344.3,p50=329458,p99=1062291,dec=2.000,hash=1.000 +macro,1000,2884.679ms,353.0(pass2 353.0),p50=2276292,p99=7292208,dec=3.000,hash=3.000 +macro,1000,2831.619ms,353.2,p50=2304750,p99=8211292,dec=3.000,hash=3.000 +``` + +The throughput win is driven by the removed redundant decodes/hashes (byte +volumes halved), not by memory effects. + +## GC / heap analysis + +### Peak heap (`getMaxMem`, process-monotonic) + +| Scenario | Pre (A) | Post (B) | Δ | +|---|---|---|---| +| micro (first scenario, cleanest) | 333.60 MB | 263.61 MB | **−21.0 %** | +| macro (whole-process peak) | 514.86 MB | 408.40 MB | **−20.7 %** | + +`gc_collections` (from `GC_getStatistics`): micro 7 / 8, macro 10 — **identical +on both branches**. Same number of GC cycles; the new code simply sits at a +lower occupied level at every point. + +### Heap-elevation series (occupied / total, MB) + +**Micro — no retention.** Occupied oscillates around a *flat* mean; `getTotalMem` +(reserved arena) is dead-flat for the entire loop on both branches. There is **no +elevation slope** in either — the transient decode/hash temporaries are allocated +and freed between the 50-msg samples and the refc arena is reused in place. + +| msg idx | A occ | A total | B occ | B total | +|---|---|---|---|---| +| 0 | 275.50 | 333.60 | 207.07 | 263.61 | +| 200 | 275.41 | 333.60 | 211.45 | 263.61 | +| 400 | 274.67 | 333.60 | 208.16 | 263.61 | +| 600 | 275.43 | 333.60 | 217.90 | 263.61 | +| 800 | 274.36 | 333.60 | 208.91 | 263.61 | +| 950 | 275.37 | 333.60 | 215.19 | 263.61 | + +- A occupied range 273.5–275.7 MB → **sawtooth amplitude ≈ 2.2 MB**. +- B occupied range 207.1–217.9 MB → **sawtooth amplitude ≈ 10.8 MB**. +- Both slopes ≈ 0 MB/100msg. Reserved total constant throughout. + +**Macro — archive retains all 1000 messages.** Occupied climbs monotonically +(this is *retention*, the archive `SortedSet` accumulating messages), not +transient churn. + +| msg idx | A occ | A total | B occ | B total | +|---|---|---|---|---| +| 50 | 308.66 | 414.61 | 235.78 | 329.06 | +| 200 | 321.78 | 414.62 | 256.83 | 329.06 | +| 400 | 353.13 | 414.64 | 283.56 | 329.07 | +| 600 | 378.20 | 414.66 | 315.62 | 408.36 | +| 800 | 405.33 | 414.85 | 336.89 | 408.38 | +| 1000 | 428.25 | 514.86 | 361.54 | 408.40 | + +- Retention slope: A ≈ **12.6 MB / 100 msg**, B ≈ **13.2 MB / 100 msg** — + statistically identical (same 1000 messages retained; payload bytes are + retained identically under value vs ref semantics). +- Reserved-arena growth: exactly **one** step each — A 414.6→514.8 MB at msg + ~650; B 329.1→408.4 MB at msg ~600. +- The new code runs a **constant ~65–75 MB lower** at every sample and peaks + 20.7 % lower. + +
Full CSV series (every 50 msgs; MB) + +``` +# scenario=micro_run1 branch=A(pre-2) idx,occMB,totMB +0,275.50,333.60 +50,275.40,333.60 +100,273.50,333.60 +150,273.95,333.60 +200,275.41,333.60 +250,275.43,333.60 +300,275.41,333.60 +350,273.58,333.60 +400,274.67,333.60 +450,273.52,333.60 +500,275.41,333.60 +550,274.95,333.60 +600,275.43,333.60 +650,275.16,333.60 +700,275.43,333.60 +750,274.01,333.60 +800,274.36,333.60 +850,273.54,333.60 +900,275.70,333.60 +950,275.37,333.60 +# scenario=macro branch=A(pre-2) idx,occMB,totMB +50,308.66,414.61 +100,313.30,414.61 +150,315.54,414.62 +200,321.78,414.62 +250,330.46,414.63 +300,345.20,414.63 +350,346.85,414.63 +400,353.13,414.64 +450,362.43,414.64 +500,376.96,414.65 +550,383.31,414.65 +600,378.20,414.66 +650,385.16,514.83 +700,395.89,514.84 +750,398.69,514.84 +800,405.33,514.85 +850,412.23,514.85 +900,423.95,514.85 +950,425.55,514.86 +1000,428.25,514.86 +# scenario=micro_run1 branch=B(post-3) idx,occMB,totMB +0,207.07,263.61 +50,215.58,263.61 +100,209.58,263.61 +150,215.05,263.61 +200,211.45,263.61 +250,207.21,263.61 +300,217.89,263.61 +350,211.46,263.61 +400,208.16,263.61 +450,212.43,263.61 +500,215.31,263.61 +550,210.43,263.61 +600,217.90,263.61 +650,211.09,263.61 +700,212.69,263.61 +750,213.75,263.61 +800,208.91,263.61 +850,207.23,263.61 +900,214.70,263.61 +950,215.19,263.61 +# scenario=macro branch=B(post-3) idx,occMB,totMB +50,235.78,329.06 +100,247.91,329.06 +150,250.84,329.06 +200,256.83,329.06 +250,261.66,329.06 +300,268.57,329.06 +350,277.31,329.07 +400,283.56,329.07 +450,290.08,329.08 +500,303.78,329.08 +550,306.02,329.08 +600,315.62,408.36 +650,315.92,408.37 +700,333.08,408.38 +750,330.20,408.38 +800,336.89,408.38 +850,346.05,408.39 +900,354.25,408.39 +950,359.48,408.40 +1000,361.54,408.40 +``` +
+ +## Verdict on the hypothesis + +> Hypothesis: retained memory ~flat (temporary copies dominate), but the OLD +> code shows a *faster heap-allocation elevation pattern* (steeper growth / +> higher sawtooth amplitude / higher peak) than the new code. + +**(a) Retained memory flat? — CONFIRMED (as a between-version statement).** +In micro there is no retention and occupied is flat on both branches. In macro +the retention *slope* is essentially identical pre vs post (~12.6 vs ~13.2 +MB/100 msg) because the same 1000 messages are retained regardless of value-vs- +ref semantics. Retained growth is a function of the workload, not the refactor. + +**(b) Allocation-elevation *slower* after the refactor? — NOT SUPPORTED by this +instrumentation.** The occupied/total sampling shows: +- Macro elevation slope is **unchanged** (retention-driven, not churn-driven). +- Micro elevation slope is **zero on both**; if anything the micro sawtooth + *amplitude* is larger post-refactor (10.8 MB vs 2.2 MB), the opposite of the + hypothesis phrasing — but this is noise-level and the reserved arena is flat. + +The reason is the refc caveat: under `--mm:refc` the eliminated transient deep +copies (async-closure env captures, Result/Option bind-outs, redundant decode +buffers) are freed **deterministically** as each temporary leaves scope, and the +arena pages are reused in place. `getTotalMem` is *peak reserved*, not +*cumulative allocated*, so it plateaus once the arena is large enough; between +two 50-msg samples the transient churn has already been allocated **and** freed. +A coarse occupied/total sawtooth therefore **cannot** observe the transient-copy +reduction — and `gc_collections` is identical (10/10), reinforcing this. + +**What the data *does* prove about memory:** the refactor lowers the **absolute +heap level** — peak heap −21 % (micro) / −20.7 % (macro), and a constant +~65–75 MB lower occupied baseline throughout the macro run. That is a real, +measurable memory win (fewer live temporaries at any instant, plus half the hash +buffers). It manifests as a **lower constant offset**, not a gentler slope or +smaller sawtooth. + +**What *would* detect the transient-copy reduction directly:** a cumulative +bytes-allocated counter (instrumenting the Nim allocator or exposing +`GC_getStatistics` cumulative fields), or malloc-level profiling (macOS +Instruments Allocations, `heaptrack`, or valgrind `massif`), or a peak-RSS probe +under memory pressure. The current harness exposes none of these — a harness +gap, not a missing win. + +## Bottom line + +| Claim | Result | +|---|---| +| Throughput up | **micro +85.9 %, macro +56.6 %** | +| Redundant decode/hash removed | hashes/msg −50 % micro, decode+hash −40–50 % macro (byte-exact) | +| Peak heap down | **−21 % / −20.7 %** | +| Retained-memory slope changed by refactor | No (retention is workload-driven, identical) | +| Old code shows steeper allocation elevation | Not observable via occupied/total sampling under refc; win is a level shift, not a slope | From c91434e509ababc91310556bb2109dd4f4301867 Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 09:57:21 +0200 Subject: [PATCH 09/10] Summary report --- docs/analysis/nocopy_summary_report.html | 358 +++++++++++++++++++++++ 1 file changed, 358 insertions(+) create mode 100644 docs/analysis/nocopy_summary_report.html diff --git a/docs/analysis/nocopy_summary_report.html b/docs/analysis/nocopy_summary_report.html new file mode 100644 index 000000000..82d7f2a32 --- /dev/null +++ b/docs/analysis/nocopy_summary_report.html @@ -0,0 +1,358 @@ + + + + + +Logos Messaging — Async Copy Elimination Report + + + + + +
+
+
λ
+

Async Copy Elimination in the Message Path

+
Deep-copy and hash-recomputation analysis of Logos Messaging under chronos / refc, + and the measured result of removing them — engineering summary report.
+
+ LOGOS MESSAGING (NIM)  ·  logos-delivery  ·  2026-07-15 +  ·  Nim 2.2.4 · --mm:refc · chronos 4.2.2 · Apple M4 +
+
+ master → experimental/chore-nocopy-e2e-perf-harness (2a4da679) + → experimental/chore-nocopy-wakumessage-refobj (3d98a28d) + → experimental/chore-nocopy-wakuenvelop (a6f1065a · be6bfed6) +
+
+
+ + +
+
+
01
+

Hypothesis & Static Analysis

+ +

Hypothesis. Under --mm:refc, chronos async gates and Result/Option idioms force + deep copies of every seq-carrying value they touch. A single inbound relayed + WakuMessage (three heap seqs + a string) was predicted to cost on the order of + 19 + F full-payload copies (F = filter peers) and 4–6 SHA-256 passes, + almost all redundant. A secondary memory hypothesis was stated up front: no real gain in retained + memory is expected (the copies are temporaries), but a slower heap-allocation elevation pattern on the + GC side. Both were put to measurement.

+ +

Where chronos deep-copies (root mechanics)

+
+ + + + + + + + + + + + + +
Copy pointMechanismrefcORC
Every {.async.} call, per seq/string paramparams lambda-lifted into the closure-iterator env (asyncmacro.nim:445-517)deep copy per paramelided (move)
complete(future, val)val: T not sink; internalValue = val (asyncfutures.nim:198-202); move() degrades to copy under refc1–2 copies1 copy
let x = await f(...)value() returns lent T; the bind to x copies1 copy1 copy
Result/Option bind-out (valueOr, ?, get)accessors are lent/templates; the let bind of a large value copies; valueOr/? also copy an lvalue Result1 copy per bindusually moved
+ +

Accumulated per-message budget on the inbound relay path (static count, verified by counters)

+
+ + + + + + + + + + + + + + + +
Cost classSitesCount / msg
Redundant proto decodes of the same byteswaku_relay/protocol.nim — ordered validator :543 · onRecv observer :271 · onValidated :326 · topicHandler :603 (+ onSend :341 outbound)4 (receiver), ≈2 avoidable copies each
Async closure env captures of the decoded messagenode/subscription_manager.nim — uniqueTopicHandler :72 + trace/filter/archive/sync/internal/legacyApp handlers :45–827 deep copies
Hash recomputation (computeMessageHash)relay :217/:553/:571/:686 · archive :101 · filter :195/:246 · store-sync :78 · node publish :1484–6 SHA-256 passes
Per-peer buffer copy (filter push)waku_filter_v2/protocol.nim:170{.async.} by-value buffer per subscribed peerF copies
Totalvs. theoretical minimum ≈ 3 (decode once · hash once · encode once)≈ 19 + F copies, 4–6 hashes
+

Phase 1 instrumentation confirmed the static analysis before any fix: the two-node + benchmark measured 6.0 decodes / 5.0 hashes per message process-wide — the receiver alone decoding the + identical proto bytes 4×. At the 10/50/150 kB payload profile this amplifies a 50 kB message + into ≈ 1 MB of heap traffic.

+
+
+ + +
+
+
02
+

Three Phases, One Branch Train

+

Each phase lives on its own local branch, chained for a PR train off master. + Every phase ends with the same harness re-run, so each claim is a measured delta, not a projection.

+
+
+
PHASE 01 — Measure first
+

E2E performance harness

+
experimental/chore-nocopy-e2e-perf-harness · 4 commits → 2a4da679
+
    +
  • Micro + macro benchmark (apps/benchmarks/message_path_bench.nim), deterministic fixed-seed workload.
  • +
  • -d:msgPathCounters: counters inside WakuMessage.decode and computeMessageHash — zero cost when undefined.
  • +
  • Baseline committed (bench_baseline.md) before touching production code.
  • +
+
Gate: counters must confirm ≥ 4 decodes and 4–6 hashes/msg — passed (6.0 / 5.0 aggregate).
+
+
+
PHASE 02 — Copies become pointers
+

WakuMessage as ref object

+
experimental/chore-nocopy-wakumessage-refobj · 5 commits → 3d98a28d
+
    +
  • Every assignment / async capture / Result bind of a message: 3-seq deep copy → pointer + refcount. No signature changes.
  • +
  • Structural ==, explicit clone(), immutable-by-convention.
  • +
  • Aliasing audit of all mutation sites; 5 real fixes (publish timestamp, ensureTimestampSet, 2× RLN proof attach, postgres nil-init). ASAN clean.
  • +
+
Gate: decode/hash counters byte-identical to baseline (no behavior change) — passed (2/2 · 6/5 unchanged).
+
+
+
PHASE 03 — Decode once, hash once
+

WakuEnvelope API break

+
experimental/chore-nocopy-wakuenvelop · 7 commits → a6f1065a
+
    +
  • WakuEnvelope = (msg · pubsubTopic · hash) as one ref; WakuRelayHandler takes the envelope; archive/filter/sync/FFI reuse envelope.hash.
  • +
  • Unused onRecv observer decode deleted; onValidated/onSend decodes gated to DEBUG.
  • +
  • Filter push buffer shared across peers (no per-peer copy).
  • +
+
Gate: receiver-side ≤ 2 decodes and exactly 1 hash per message — passed (2.0 / 1.0).
+
+
+
+
+ + +
+
+
03
+

Measurement Method

+
+ + + + + + + + + + + + + +
ScenarioEntry endpointExit endpoint (measured)What it isolates
micro — in-process, no networkraw proto bytes into the relay’s registered ordered validator, then its topicHandlerreturn of the full dispatch chain (trace → archive insert → store-sync → internal event)the receiver pipeline, low-noise (0.4–4 % variance)
macro — end-to-endnodeA.publish() — public node API → gossipsub → loopback TCPapplication-level relay handler on node B, firing only after validation, decode, dispatch and archive insert completethe full pipeline incl. libp2p observers
+

Workload: fixed seed 42, payload mix 10 / 50 / 150 kB at 25 / 50 / 25 %, N = 1000 + 100 warmup, + unique payload prefix + timestamp (no gossipsub dedup). Publisher flow-controlled to a 32-message window. + Decode/hash counters live inside the codec and hash procs themselves — they count reality, not expectations. + Peak heap via getMaxMem(); heap series sampled every 50 messages; retained delta via + getOccupiedMem() after GC_fullCollect().

+

Final comparison ran both branch heads back-to-back in one session on the same machine + (one post-refactor pass discarded for a >5 % variance load spike, rerun within bounds). Not covered by + these numbers: RLN validation, filter push to remote peers, REST/FFI boundary, real network latency.

+
+
+ + +
+
+
04
+

Measured Gains — Before vs. After

+

“Before” = pre-Phase-2 head (2a4da679: original production code + harness). + “After” = post-Phase-3 head (a6f1065a). Same-session A/B.

+ +
+
+86 %
micro throughput
1,223 → 2,274 msg/s
+
+57 %
e2e (macro) throughput
225 → 353 msg/s
+
−21 %
peak heap (both scenarios)
515 → 408 MB macro
+
4 → 2
receiver decodes / msg
hashes: receiver 3 → 1
+
+ +
+ Before (pre-Phase-2) + After (post-Phase-3) + bars are supplementary — full values in the tables below +
+ +
+
Throughput (msg/s, higher is better)
+
+
micro — before
1,223
+
micro — after
2,274
+
+
+
macro e2e — before
225
+
macro e2e — after
353
+
+
+ +
+
Redundant work per message (macro, process aggregate — lower is better)
+
+
proto decodes — before
6.0
+
proto decodes — after
3.0
+
+
+
SHA-256 passes — before
5.0
+
SHA-256 passes — after
3.0
+
+
+
bytes hashed — before
325 MB / 1000 msgs
+
bytes hashed — after
195 MB / 1000 msgs
+
+
+ +

Full comparison

+
+ + + + + + + + + + + + +
MetricBeforeAfterΔ
micro msg/s1,222.82,273.7+85.9 %
micro p50 / p99 per msg630 µs / 1.94 ms338 µs / 1.09 ms−46 % / −44 %
micro decodes / hashes per msg2.0 / 2.02.0 / 1.0hash −50 %
macro msg/s (e2e)225.5353.1+56.6 %
macro p50 / p99 inter-arrival3.55 ms / 10.7 ms2.29 ms / 7.5 ms−35 % / −30 %
macro decodes / hashes per msg (aggregate)6.0 / 5.03.0 / 3.0−50 % / −40 %
  — receiver-side only4 dec / 3 hash2 dec / 1 hashgate met
bytes decoded / hashed per 1000 msgs390 / 325 MB195 / 195 MB−50 % / −40 %
peak heap getMaxMem (micro / macro)333.6 / 514.9 MB263.6 / 408.4 MB−21 % both
GC collections (micro / macro)7–8 / 107–8 / 10identical
retained-mem slope, macro≈ 12.6 MB / 100 msg≈ 13.2 MB / 100 msgflat (archive-driven)
+ +

Memory hypothesis — verdict

+
+
+
CONFIRMED
+

(a) No real retained-memory gain

+

Retention slope is identical before and after — it is the archive holding the same 1000 messages. + The eliminated copies were temporaries, exactly as hypothesized.

+
+
+
REFUTED — LEVEL SHIFT INSTEAD
+

(b) Slower heap-elevation pattern

+

Slope and collection counts are unchanged: under refc the transient copies were refcount-freed + deterministically and never accumulated toward collection triggers. The win manifests as a + −21 % peak heap and a ~65–75 MB lower level throughout the run — the copies cost CPU + (memcpy + alloc/free churn) and peak footprint, not GC-cycle pressure. Quantifying raw churn would + require cumulative-allocation counters or malloc profiling.

+
+
+
+
+ + +
+
+

Conclusions & Open Items

+

The message path now performs one decode per trust boundary and one hash per message, + carried by reference through every async gate. Verification at each phase: representative suites green + (~160+ tests), ASAN clean on the dispatch path, counters byte-exact where no change was claimed.

+
    +
  • Per-shard payload byte gauges (waku_relay_*_msg_bytes_per_shard) now populate only at DEBUG/TRACE — decide before merge if dashboards need them.
  • +
  • Archive/filter keep thin (topic, msg) compat overloads; removable later.
  • +
  • Lightpush validateMessage double-encode deferred (needs PushMessageHandler signature change).
  • +
  • Under a future ORC build the remaining complete()/await-bind copies persist — returning refs from async procs stays best practice.
  • +
+

Sources: docs/analysis/async_copy_analysis.md · async_copy_fix_plan.md · plan_phase{1,2,3}_*.md · + bench_baseline.md · phase2_mutation_audit.md · speed_gain_pre2_post3.md — all committed on the branch train. Local branches only; nothing pushed.

+
Logos.coλ Logos Messaging · Engineering Report · 2026-07-15
+
+
+ + + From a35b1a8a1ae04e7633a528828da589fb96ef4b58 Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Wed, 15 Jul 2026 10:13:45 +0200 Subject: [PATCH 10/10] docs: markdown version of the nocopy summary report GitHub-renderable mirror of nocopy_summary_report.html. Co-Authored-By: Claude Fable 5 --- docs/analysis/nocopy_summary_report.md | 163 +++++++++++++++++++++++++ 1 file changed, 163 insertions(+) create mode 100644 docs/analysis/nocopy_summary_report.md diff --git a/docs/analysis/nocopy_summary_report.md b/docs/analysis/nocopy_summary_report.md new file mode 100644 index 000000000..4b0d6a279 --- /dev/null +++ b/docs/analysis/nocopy_summary_report.md @@ -0,0 +1,163 @@ +# λ Async Copy Elimination in the Message Path — Summary Report + +Deep-copy and hash-recomputation analysis of Logos Messaging under chronos / refc, and the +measured result of removing them. + +> **Logos Messaging (Nim) · logos-delivery · 2026-07-15** +> Nim 2.2.4 · `--mm:refc` · chronos 4.2.2 · Apple M4 +> +> Branch train: `master` → `experimental/chore-nocopy-e2e-perf-harness` (2a4da679) +> → `experimental/chore-nocopy-wakumessage-refobj` (3d98a28d) +> → `experimental/chore-nocopy-wakuenvelop` (a6f1065a · be6bfed6) +> +> HTML version: [`nocopy_summary_report.html`](nocopy_summary_report.html) + +--- + +## 01 — Hypothesis & Static Analysis + +**Hypothesis.** Under `--mm:refc`, chronos async gates and Result/Option idioms force deep copies +of every `seq`-carrying value they touch. A single inbound relayed `WakuMessage` (three heap seqs + +a string) was predicted to cost on the order of **19 + F full-payload copies** (F = filter peers) +and **4–6 SHA-256 passes**, almost all redundant. A secondary memory hypothesis was stated up +front: *no real gain in retained memory is expected (the copies are temporaries), but a slower +heap-allocation elevation pattern on the GC side.* Both were put to measurement. + +### Where chronos deep-copies (root mechanics) + +| Copy point | Mechanism | refc | ORC | +|---|---|---|---| +| Every `{.async.}` call, per seq/string param | params lambda-lifted into the closure-iterator env (`asyncmacro.nim:445-517`) | deep copy per param | elided (move) | +| `complete(future, val)` | `val: T` not `sink`; `internalValue = val` (`asyncfutures.nim:198-202`); `move()` degrades to copy under refc | 1–2 copies | 1 copy | +| `let x = await f(...)` | `value()` returns `lent T`; the bind to `x` copies | 1 copy | 1 copy | +| Result/Option bind-out (`valueOr`, `?`, `get`) | accessors are `lent`/templates; the `let` bind of a large value copies; `valueOr`/`?` also copy an lvalue Result | 1 copy per bind | usually moved | + +### Accumulated per-message budget on the inbound relay path (static count, verified by counters) + +| Cost class | Sites | Count / msg | +|---|---|---| +| **Redundant proto decodes** of the same bytes | `waku_relay/protocol.nim` — ordered validator :543 · onRecv observer :271 · onValidated :326 · topicHandler :603 (+ onSend :341 outbound) | 4 (receiver), ≈2 avoidable copies each | +| **Async closure env captures** of the decoded message | `node/subscription_manager.nim` — uniqueTopicHandler :72 + trace/filter/archive/sync/internal/legacyApp handlers :45–82 | 7 deep copies | +| **Hash recomputation** (`computeMessageHash`) | relay :217/:553/:571/:686 · archive :101 · filter :195/:246 · store-sync :78 · node publish :148 | 4–6 SHA-256 passes | +| **Per-peer buffer copy** (filter push) | `waku_filter_v2/protocol.nim:170` — `{.async.}` by-value buffer per subscribed peer | F copies | +| **Total** | vs. theoretical minimum ≈ 3 (decode once · hash once · encode once) | **≈ 19 + F copies, 4–6 hashes** | + +Phase 1 instrumentation confirmed the static analysis before any fix: the two-node benchmark +measured 6.0 decodes / 5.0 hashes per message process-wide — the receiver alone decoding the +identical proto bytes 4×. At the 10/50/150 kB payload profile this amplifies a 50 kB message into +≈ 1 MB of heap traffic. + +--- + +## 02 — Three Phases, One Branch Train + +Each phase lives on its own local branch, chained for a PR train off `master`. Every phase ends +with the same harness re-run, so each claim is a measured delta, not a projection. + +### Phase 01 — Measure first: e2e performance harness + +`experimental/chore-nocopy-e2e-perf-harness` · 4 commits → `2a4da679` + +- Micro + macro benchmark (`apps/benchmarks/message_path_bench.nim`), deterministic fixed-seed workload. +- `-d:msgPathCounters`: counters inside `WakuMessage.decode` and `computeMessageHash` — zero cost when undefined. +- Baseline committed (`bench_baseline.md`) before touching production code. + +> **Gate:** counters must confirm ≥ 4 decodes and 4–6 hashes/msg — **passed** (6.0 / 5.0 aggregate). + +### Phase 02 — Copies become pointers: WakuMessage as `ref object` + +`experimental/chore-nocopy-wakumessage-refobj` · 5 commits → `3d98a28d` + +- Every assignment / async capture / Result bind of a message: 3-seq deep copy → pointer + refcount. No signature changes. +- Structural `==`, explicit `clone()`, immutable-by-convention. +- Aliasing audit of all mutation sites; **5 real fixes** (publish timestamp, `ensureTimestampSet`, 2× RLN proof attach, postgres nil-init). ASAN clean. + +> **Gate:** decode/hash counters byte-identical to baseline (no behavior change) — **passed** (2/2 · 6/5 unchanged). + +### Phase 03 — Decode once, hash once: WakuEnvelope API break + +`experimental/chore-nocopy-wakuenvelop` · 7 commits → `a6f1065a` + +- `WakuEnvelope` = (msg · pubsubTopic · hash) as one ref; `WakuRelayHandler` takes the envelope; archive/filter/sync/FFI reuse `envelope.hash`. +- Unused onRecv observer decode deleted; onValidated/onSend decodes gated to DEBUG. +- Filter push buffer shared across peers (no per-peer copy). + +> **Gate:** receiver-side ≤ 2 decodes and exactly 1 hash per message — **passed** (2.0 / 1.0). + +--- + +## 03 — Measurement Method + +| Scenario | Entry endpoint | Exit endpoint (measured) | What it isolates | +|---|---|---|---| +| **micro** — in-process, no network | raw proto bytes into the relay's registered *ordered validator*, then its *topicHandler* | return of the full dispatch chain (trace → archive insert → store-sync → internal event) | the receiver pipeline, low-noise (0.4–4 % variance) | +| **macro** — end-to-end | `nodeA.publish()` — public node API → gossipsub → loopback TCP | application-level relay handler on node B, firing only after validation, decode, dispatch and archive insert complete | the full pipeline incl. libp2p observers | + +Workload: fixed seed 42, payload mix **10 / 50 / 150 kB at 25 / 50 / 25 %**, N = 1000 + 100 warmup, +unique payload prefix + timestamp (no gossipsub dedup). Publisher flow-controlled to a 32-message +window. Decode/hash counters live inside the codec and hash procs themselves — they count reality, +not expectations. Peak heap via `getMaxMem()`; heap series sampled every 50 messages; retained +delta via `getOccupiedMem()` after `GC_fullCollect()`. + +Final comparison ran both branch heads back-to-back in one session on the same machine (one +post-refactor pass discarded for a >5 % variance load spike, rerun within bounds). Not covered by +these numbers: RLN validation, filter push to remote peers, REST/FFI boundary, real network latency. + +--- + +## 04 — Measured Gains — Before vs. After + +"Before" = pre-Phase-2 head (`2a4da679`: original production code + harness). +"After" = post-Phase-3 head (`a6f1065a`). Same-session A/B. + +| | | | +|---|---|---| +| **+86 %** micro throughput (1,223 → 2,274 msg/s) | **+57 %** e2e throughput (225 → 353 msg/s) | **−21 %** peak heap (515 → 408 MB macro) | + +### Full comparison + +| Metric | Before | After | Δ | +|---|---:|---:|---:| +| micro msg/s | 1,222.8 | 2,273.7 | **+85.9 %** | +| micro p50 / p99 per msg | 630 µs / 1.94 ms | 338 µs / 1.09 ms | −46 % / −44 % | +| micro decodes / hashes per msg | 2.0 / 2.0 | 2.0 / 1.0 | hash −50 % | +| macro msg/s (e2e) | 225.5 | 353.1 | **+56.6 %** | +| macro p50 / p99 inter-arrival | 3.55 ms / 10.7 ms | 2.29 ms / 7.5 ms | −35 % / −30 % | +| macro decodes / hashes per msg (aggregate) | 6.0 / 5.0 | 3.0 / 3.0 | −50 % / −40 % | +| — receiver-side only | 4 dec / 3 hash | **2 dec / 1 hash** | gate met | +| bytes decoded / hashed per 1000 msgs | 390 / 325 MB | 195 / 195 MB | −50 % / −40 % | +| peak heap `getMaxMem` (micro / macro) | 333.6 / 514.9 MB | 263.6 / 408.4 MB | **−21 % both** | +| GC collections (micro / macro) | 7–8 / 10 | 7–8 / 10 | identical | +| retained-mem slope, macro | ≈ 12.6 MB / 100 msg | ≈ 13.2 MB / 100 msg | flat (archive-driven) | + +### Memory hypothesis — verdict + +**(a) No real retained-memory gain — CONFIRMED.** Retention slope is identical before and after — +it is the archive holding the same 1000 messages. The eliminated copies were temporaries, exactly +as hypothesized. + +**(b) Slower heap-elevation pattern — REFUTED; level shift instead.** Slope and collection counts +are unchanged: under refc the transient copies were refcount-freed deterministically and never +accumulated toward collection triggers. The win manifests as a **−21 % peak heap** and a +~65–75 MB lower level throughout the run — the copies cost CPU (memcpy + alloc/free churn) and +peak footprint, not GC-cycle pressure. Quantifying raw churn would require cumulative-allocation +counters or malloc profiling. + +--- + +## Conclusions & Open Items + +The message path now performs one decode per trust boundary and one hash per message, carried by +reference through every async gate. Verification at each phase: representative suites green +(~160+ tests), ASAN clean on the dispatch path, counters byte-exact where no change was claimed. + +- Per-shard payload byte gauges (`waku_relay_*_msg_bytes_per_shard`) now populate only at DEBUG/TRACE — decide before merge if dashboards need them. +- Archive/filter keep thin `(topic, msg)` compat overloads; removable later. +- Lightpush `validateMessage` double-encode deferred (needs `PushMessageHandler` signature change). +- Under a future ORC build the remaining `complete()`/await-bind copies persist — returning refs from async procs stays best practice. + +Sources: `async_copy_analysis.md` · `async_copy_fix_plan.md` · `plan_phase{1,2,3}_*.md` · +`bench_baseline.md` · `phase2_mutation_audit.md` · `speed_gain_pre2_post3.md` — all committed on +the branch train. + +*Logos.co · λ Logos Messaging · Engineering Report · 2026-07-15*