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.
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`).
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 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 (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)
+
+
+
+
+
+
Redundant work per message (macro, process aggregate — lower is better)
+
+
+
+
+
+
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
+
+
+
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.
+
+
+
+
+
+
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*
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 |
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/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()
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_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/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",
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.
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/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_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
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: