test(e2e): port remaining wrapper tests from the interop repo (#4077)

* test(e2e): port remaining wrapper tests from the interop repo

The wrapper suite that moved into tests-e2e (#4027) was a reworked subset of
the one still living in logos-delivery-interop-tests. Comparing both sides
showed 21 tests here against 46 there, with no overlap in the delta: the 25
missing tests cover scenarios the reworked set never included.

Ports those 25 tests, bringing the in-repo suite to the full 46:
  - 12 send scenarios: s01 (nil/destroyed handle), s03, s04, s05, s11, s13,
    s16, s18 (both orderings), s25, s29
  - 7 channel lifecycle tests (rc01-rc04)
  - 6 wrapper corner cases: auto port allocation, MyBoundPorts, ENR

Supporting changes the ported tests need:
  - wrapper_helpers: get_node_tcp_port, get_node_bound_ports, enr_udp_port
  - WrapperManager: channel_create/send/close, destroy_keep_ctx
  - vendored binding refreshed to the revision exposing the channel API
    (additive only; cffi resolves symbols lazily, so nothing existing moves)

Two Edge senders were fixed while porting. build_node_config defaults relay
and store to True, and the flat-JSON config path applies mode=Edge before
explicit fields, so those defaults win: the Edge nodes in s11/s16/s25 came up
as relay and store servers and exercised the relay path instead of lightpush.
They now set relay=False and store=False, matching test_send_e2e_part2. s16
also dropped lightpush=True, which mounts the lightpush server and fails node
start once relay is off; the lightpush client mounts unconditionally.

Suite goes from 21 to 46 functions (53 collected). The docker subset grows
from 3 to 5 as s11 and s25 need a store peer. Local run against a freshly
built library: 45 passed, 2 skipped, 1 xfailed.

* test(e2e): enable autosharding in the channel lifecycle tests

channel_create subscribes to the channel's content topic since #4081, and
resolving that topic to a shard needs autosharding. build_node_config leaves
numShardsInNetwork at 0 and cluster 198 has no preset, so these nodes came up
with static sharding and every channel_create failed with "autosharding is not
configured; pass an explicit shard".

Adds numShardsInNetwork=1 to the six tests that create a channel, matching what
every other wrapper test that touches the send or channel API already does.
rc02 is left alone: channel_send rejects on the id lookup before any shard is
resolved.

Verified locally against a fresh build: the five tests that complete now pass
and the error string is gone from the run.
This commit is contained in:
Egor Rachkovskii 2026-08-10 08:36:28 +01:00 committed by GitHub
parent 4a85db1b6a
commit 342a965370
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
12 changed files with 1717 additions and 2 deletions

1
.gitignore vendored
View File

@ -91,6 +91,7 @@ nimbledeps
# Python bytecode from tests/simulator
__pycache__/
*.pyc
.venv*/
# Emitted by genBindings() during the liblogosdelivery build.
library/generated/

View File

@ -21,8 +21,8 @@ pip install -r tests-e2e/requirements.txt
# 3. Run (from tests-e2e/)
cd tests-e2e
pytest tests/wrappers_tests -m "not docker_required" # 18 pure-binding tests
pytest tests/wrappers_tests -m docker_required # 3 tests that also need a Docker nwaku peer (S19/S20/S31)
pytest tests/wrappers_tests -m "not docker_required" # 48 pure-binding tests
pytest tests/wrappers_tests -m docker_required # 5 tests that also need a Docker nwaku peer (S11/S19/S20/S25/S31)
```
## CI

View File

@ -1,6 +1,8 @@
from __future__ import annotations
import base64
import json
import re
import threading
import time
from typing import Optional
@ -195,6 +197,53 @@ def get_node_multiaddr(node) -> str:
return addr
# Matches the /tcp/<port>/ segment in a libp2p multiaddr.
TCP_PORT_RE = re.compile(r"/tcp/(\d+)/")
def get_node_tcp_port(node) -> int:
"""Return the TCP port the node advertises in its multiaddr."""
multiaddr = get_node_multiaddr(node)
match = TCP_PORT_RE.search(multiaddr)
if not match:
raise RuntimeError(f"multiaddr missing /tcp/<port>/ segment: {multiaddr!r}")
return int(match.group(1))
def get_node_bound_ports(node) -> dict:
"""Return the MyBoundPorts debug info .
Keys: tcp, webSocket, quic, rest, discv5Udp, metrics. A value of 0 means the
service is disabled or did not bind.
"""
result = node.get_node_info("MyBoundPorts")
if result.is_err():
raise RuntimeError(f"MyBoundPorts query failed: {result.err()}")
return json.loads(result.ok_value)
def enr_udp_port(enr_uri: str) -> int:
"""Extract the advertised udp port from a text-encoded ENR.
An ENR is "enr:" + base64url(RLP). Instead of pulling in a full RLP
decoder, find the "udp" key in the raw bytes and read the value after it.
"""
if not enr_uri.startswith("enr:"):
raise RuntimeError(f"not an ENR URI: {enr_uri!r}")
b64 = enr_uri[len("enr:") :]
payload = base64.urlsafe_b64decode(b64 + "=" * (-len(b64) % 4))
key = payload.find(b"\x83udp") # "udp" encoded as a 3-byte RLP string
if key == -1:
raise RuntimeError(f"ENR has no udp entry: {enr_uri!r}")
prefix = payload[key + 4]
if prefix < 0x80: # values < 128 are encoded as a single byte
return prefix
size = prefix - 0x80 # short string: prefix is 0x80 + length
return int.from_bytes(payload[key + 5 : key + 5 + size], "big")
def create_message_bindings(**overrides) -> dict:
envelope = {
"contentTopic": DEFAULT_CONTENT_TOPIC,

View File

@ -62,6 +62,10 @@ class WrapperManager:
def stop_and_destroy(self, *, timeout_s: float = 20.0) -> Result[int, str]:
return self._node.stop_and_destroy(timeout_s=timeout_s)
def destroy_keep_ctx(self, *, timeout_s: float = 20.0) -> Result[int, str]:
"""Pass-through for NodeWrapper.destroy_keep_ctx — see that method."""
return self._node.destroy_keep_ctx(timeout_s=timeout_s)
def subscribe_content_topic(self, content_topic: str, *, timeout_s: float = 20.0) -> Result[int, str]:
return self._node.subscribe_content_topic(content_topic, timeout_s=timeout_s)
@ -71,6 +75,22 @@ class WrapperManager:
def send_message(self, message: dict, *, timeout_s: float = 20.0) -> Result[str, str]:
return self._node.send_message(message, timeout_s=timeout_s)
def channel_create(
self,
channel_id: str,
content_topic: str,
sender_id: str,
*,
timeout_s: float = 20.0,
) -> Result[str, str]:
return self._node.channel_create(channel_id, content_topic, sender_id, timeout_s=timeout_s)
def channel_send(self, channel_id: str, message: dict, *, timeout_s: float = 20.0) -> Result[str, str]:
return self._node.channel_send(channel_id, message, timeout_s=timeout_s)
def channel_close(self, channel_id: str, *, timeout_s: float = 20.0) -> Result[str, str]:
return self._node.channel_close(channel_id, timeout_s=timeout_s)
def get_available_node_info_ids(self, *, timeout_s: float = 20.0) -> Result[list[str], str]:
return self._node.get_available_node_info_ids(timeout_s=timeout_s)

View File

@ -0,0 +1,148 @@
"""S01 helpers: invoke send() against an invalid handle in an isolated process.
Run via:
python -m tests.wrappers_tests.helpers.send_invalid_handle nil <marker>
python -m tests.wrappers_tests.helpers.send_invalid_handle destroyed <marker>
Prints a single line to stdout starting with <marker>, followed by a JSON
payload describing the outcome. Runs in its own process so that a missing
C-ABI guard (which can SIGSEGV) fails the parent test cleanly instead of
taking the pytest runner down.
Cases:
- nil: send() on a wrapper built with ctx=ffi.NULL.
- destroyed: send() after destroy_keep_ctx() self.ctx stays non-nil so the
call reaches the C side with the original (now-stale) pointer.
"""
import json
import sys
from pathlib import Path
def _ensure_bindings_on_path() -> None:
# The wrapper module lives outside the project tree, under vendor/.
# Helper file is at <root>/tests/wrappers_tests/helpers/<this>.py.
project_root = Path(__file__).resolve().parents[3]
bindings_path = project_root / "vendor" / "logos-delivery-python-bindings" / "waku"
if str(bindings_path) not in sys.path:
sys.path.insert(0, str(bindings_path))
if str(project_root) not in sys.path:
sys.path.insert(0, str(project_root))
def _emit(marker: str, payload: dict) -> None:
print(marker + json.dumps(payload))
def _run_nil_handle(marker: str) -> None:
from wrapper import NodeWrapper, ffi # type: ignore[import-not-found]
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import create_message_bindings
sender = WrapperManager(NodeWrapper(ctx=ffi.NULL, config_buffer=None, event_cb_handler=None))
send_result = sender.send_message(message=create_message_bindings())
_emit(
marker,
{
"is_ok": send_result.is_ok(),
"ok": send_result.ok_value if send_result.is_ok() else None,
"err": send_result.err() if send_result.is_err() else None,
},
)
def _run_destroyed_handle(marker: str) -> None:
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import EventCollector, create_message_bindings
from tests.wrappers_tests.conftest import build_node_config
collector = EventCollector()
create_result = WrapperManager.create_and_start(
config=build_node_config(),
event_cb=collector.event_callback,
)
if create_result.is_err():
_emit(
marker,
{
"stage": "create_and_start",
"is_ok": False,
"ok": None,
"err": create_result.err(),
"events_after_send": [],
},
)
return
sender = create_result.ok_value
stop_result = sender.stop_node()
if stop_result.is_err():
_emit(
marker,
{
"stage": "stop_node",
"is_ok": False,
"ok": None,
"err": stop_result.err(),
"events_after_send": [],
},
)
return
destroy_result = sender.destroy_keep_ctx()
if destroy_result.is_err():
_emit(
marker,
{
"stage": "destroy_keep_ctx",
"is_ok": False,
"ok": None,
"err": destroy_result.err(),
"events_after_send": [],
},
)
return
events_before_send = len(collector.events)
send_result = sender.send_message(message=create_message_bindings())
new_events = collector.events[events_before_send:]
_emit(
marker,
{
"stage": "send_message",
"is_ok": send_result.is_ok(),
"ok": send_result.ok_value if send_result.is_ok() else None,
"err": send_result.err() if send_result.is_err() else None,
"events_after_send": [str(e) for e in new_events],
},
)
CASES = {
"nil": _run_nil_handle,
"destroyed": _run_destroyed_handle,
}
def main() -> int:
if len(sys.argv) != 3 or sys.argv[1] not in CASES:
cases = "|".join(CASES)
print(f"usage: send_invalid_handle <{cases}> <result_marker>", file=sys.stderr)
return 2
case, marker = sys.argv[1], sys.argv[2]
_ensure_bindings_on_path()
CASES[case](marker)
return 0
if __name__ == "__main__":
sys.exit(main())

View File

@ -0,0 +1,228 @@
import uuid
import pytest
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import EventCollector, create_message_bindings
from src.libs.common import delay
CHANNEL_ID = "rc01-channel"
CONTENT_TOPIC = "/test/1/channel/proto"
SENDER_ID = "rc01-sender"
UNKNOWN_CHANNEL_ID = "rc02-unknown-channel"
RC03_CHANNEL_ID = "rc03-channel"
RC04_CHANNEL_ID = "rc04-channel"
RC04_EPHEMERAL_CHANNEL_ID = "rc04-ephemeral-channel"
RC04_DISTINCT_CHANNEL_ID = "rc04-distinct-channel"
# Give the reliable channel manager a moment to settle a freshly-created
# channel before we act on it.
CHANNEL_SETTLE_S = 1
@pytest.mark.smoke
class TestChannelLifecycle:
def test_rc01_create_channel_duplicate_rejected(self, node_config):
"""RC01: create a channel; a duplicate create with the same id is rejected.
The reliable channel manager is mounted automatically for a Core-mode
node, so no store / --reliability is needed. No events are expected.
"""
# channel_create subscribes to its content topic, which needs autosharding.
node_config.update({"mode": "Core", "numShardsInNetwork": 1})
collector = EventCollector()
create_node_result = WrapperManager.create_and_start(config=node_config, event_cb=collector.event_callback)
assert create_node_result.is_ok(), f"Failed to create and start node: {create_node_result.err()}"
node = create_node_result.ok_value
try:
create_result = node.channel_create(CHANNEL_ID, CONTENT_TOPIC, SENDER_ID)
assert create_result.is_ok(), f"channel_create failed: {create_result.err()}"
assert create_result.ok_value == CHANNEL_ID, f"channel_create returned unexpected id: {create_result.ok_value!r}"
duplicate_result = node.channel_create(CHANNEL_ID, CONTENT_TOPIC, SENDER_ID)
assert duplicate_result.is_err(), f"duplicate channel_create must fail, got Ok({duplicate_result.ok_value!r})"
assert f"channel already exists: {CHANNEL_ID}" in duplicate_result.err(), f"unexpected error message: {duplicate_result.err()!r}"
assert collector.events == [], f"expected no events, got: {collector.events}"
finally:
node.stop_and_destroy()
def test_rc02_send_on_unknown_channel_rejected(self, node_config):
"""RC02: channel_send() to an id that was never created is rejected.
The send must Err with "unknown channel: <id>" and emit no events.
"""
node_config.update({"mode": "Core"})
collector = EventCollector()
create_node_result = WrapperManager.create_and_start(config=node_config, event_cb=collector.event_callback)
assert create_node_result.is_ok(), f"Failed to create and start node: {create_node_result.err()}"
node = create_node_result.ok_value
try:
send_result = node.channel_send(UNKNOWN_CHANNEL_ID, create_message_bindings())
assert send_result.is_err(), f"channel_send on unknown channel must fail, got Ok({send_result.ok_value!r})"
assert f"unknown channel: {UNKNOWN_CHANNEL_ID}" in send_result.err(), f"unexpected error message: {send_result.err()!r}"
assert collector.events == [], f"expected no events, got: {collector.events}"
finally:
node.stop_and_destroy()
def test_rc03_close_channel_close_unknown_rejected(self, node_config):
"""RC03: close an existing channel; closing again / an unknown id is rejected.
First close succeeds; a second close of the same (now-gone) id and a
close of a never-created id both Err with "unknown channel: <id>".
Lifecycle teardown + idempotency guard. No events are expected.
"""
node_config.update({"mode": "Core", "numShardsInNetwork": 1})
collector = EventCollector()
create_node_result = WrapperManager.create_and_start(config=node_config, event_cb=collector.event_callback)
assert create_node_result.is_ok(), f"Failed to create and start node: {create_node_result.err()}"
node = create_node_result.ok_value
try:
create_result = node.channel_create(RC03_CHANNEL_ID, CONTENT_TOPIC, SENDER_ID)
assert create_result.is_ok(), f"channel_create failed: {create_result.err()}"
delay(CHANNEL_SETTLE_S)
close_result = node.channel_close(RC03_CHANNEL_ID)
assert close_result.is_ok(), f"channel_close on an existing channel failed: {close_result.err()}"
reclose_result = node.channel_close(RC03_CHANNEL_ID)
assert reclose_result.is_err(), f"second channel_close must fail, got Ok({reclose_result.ok_value!r})"
assert f"unknown channel: {RC03_CHANNEL_ID}" in reclose_result.err(), f"unexpected error message: {reclose_result.err()!r}"
unknown_close_result = node.channel_close(UNKNOWN_CHANNEL_ID)
assert unknown_close_result.is_err(), f"channel_close on unknown channel must fail, got Ok({unknown_close_result.ok_value!r})"
assert f"unknown channel: {UNKNOWN_CHANNEL_ID}" in unknown_close_result.err(), f"unexpected error message: {unknown_close_result.err()!r}"
assert collector.events == [], f"expected no events, got: {collector.events}"
finally:
node.stop_and_destroy()
def test_rc03_close_random_id_leaves_open_channel_intact(self, node_config):
"""RC03 variant: closing a random unknown id while a real channel is open.
A genuine channel is created and left open; closing a random id Errs with
"unknown channel: <id>" and must not disturb the open channel, which then
still closes cleanly. No events are expected.
"""
node_config.update({"mode": "Core", "numShardsInNetwork": 1})
channel_id = f"rc03-open-{uuid.uuid4()}"
random_id = f"rc03-random-{uuid.uuid4()}"
collector = EventCollector()
create_node_result = WrapperManager.create_and_start(config=node_config, event_cb=collector.event_callback)
assert create_node_result.is_ok(), f"Failed to create and start node: {create_node_result.err()}"
node = create_node_result.ok_value
try:
create_result = node.channel_create(channel_id, CONTENT_TOPIC, SENDER_ID)
assert create_result.is_ok(), f"channel_create failed: {create_result.err()}"
delay(CHANNEL_SETTLE_S)
random_close_result = node.channel_close(random_id)
assert random_close_result.is_err(), f"channel_close on a random id must fail, got Ok({random_close_result.ok_value!r})"
assert f"unknown channel: {random_id}" in random_close_result.err(), f"unexpected error message: {random_close_result.err()!r}"
# The real channel was left untouched and still closes cleanly.
close_result = node.channel_close(channel_id)
assert close_result.is_ok(), f"open channel_close failed after a bad-id close: {close_result.err()}"
assert collector.events == [], f"expected no events, got: {collector.events}"
finally:
node.stop_and_destroy()
def test_rc04_empty_payload_rejected_and_send_returns_handle(self, node_config):
"""RC04: empty payload is rejected; a non-empty send returns a handle synchronously.
channel_send with an empty payload Errs with "empty payload"; a normal
channel_send returns Ok(channelReqId) immediately (delivery is async).
Input validation + the immediate-handle contract of send.
"""
node_config.update({"mode": "Core", "numShardsInNetwork": 1})
create_node_result = WrapperManager.create_and_start(config=node_config)
assert create_node_result.is_ok(), f"Failed to create and start node: {create_node_result.err()}"
node = create_node_result.ok_value
try:
create_result = node.channel_create(RC04_CHANNEL_ID, CONTENT_TOPIC, SENDER_ID)
assert create_result.is_ok(), f"channel_create failed: {create_result.err()}"
delay(CHANNEL_SETTLE_S)
empty_result = node.channel_send(RC04_CHANNEL_ID, create_message_bindings(payload=""))
assert empty_result.is_err(), f"channel_send with an empty payload must fail, got Ok({empty_result.ok_value!r})"
assert "empty payload" in empty_result.err(), f"unexpected error message: {empty_result.err()!r}"
send_result = node.channel_send(RC04_CHANNEL_ID, create_message_bindings())
assert send_result.is_ok(), f"channel_send with a non-empty payload failed: {send_result.err()}"
assert send_result.ok_value, f"channel_send must return a non-empty channelReqId handle, got: {send_result.ok_value!r}"
finally:
node.stop_and_destroy()
def test_rc04_ephemeral_send_returns_handle(self, node_config):
"""RC04 variant: an ephemeral channel_send still returns a handle synchronously.
The ephemeral flag only affects downstream persistence / terminal events,
not the send contract, so channel_send returns Ok(channelReqId) immediately.
"""
node_config.update({"mode": "Core", "numShardsInNetwork": 1})
create_node_result = WrapperManager.create_and_start(config=node_config)
assert create_node_result.is_ok(), f"Failed to create and start node: {create_node_result.err()}"
node = create_node_result.ok_value
try:
create_result = node.channel_create(RC04_EPHEMERAL_CHANNEL_ID, CONTENT_TOPIC, SENDER_ID)
assert create_result.is_ok(), f"channel_create failed: {create_result.err()}"
delay(CHANNEL_SETTLE_S)
send_result = node.channel_send(RC04_EPHEMERAL_CHANNEL_ID, create_message_bindings(ephemeral=True))
assert send_result.is_ok(), f"ephemeral channel_send failed: {send_result.err()}"
assert send_result.ok_value, f"channel_send must return a non-empty channelReqId handle, got: {send_result.ok_value!r}"
finally:
node.stop_and_destroy()
def test_rc04_sequential_sends_return_distinct_handles(self, node_config):
"""RC04 variant: back-to-back sends each return a distinct handle.
Every channel_send is its own request, so two sends on the same channel
return two different (non-empty) channelReqId handles.
"""
node_config.update({"mode": "Core", "numShardsInNetwork": 1})
create_node_result = WrapperManager.create_and_start(config=node_config)
assert create_node_result.is_ok(), f"Failed to create and start node: {create_node_result.err()}"
node = create_node_result.ok_value
try:
create_result = node.channel_create(RC04_DISTINCT_CHANNEL_ID, CONTENT_TOPIC, SENDER_ID)
assert create_result.is_ok(), f"channel_create failed: {create_result.err()}"
delay(CHANNEL_SETTLE_S)
first_result = node.channel_send(RC04_DISTINCT_CHANNEL_ID, create_message_bindings())
assert first_result.is_ok(), f"first channel_send failed: {first_result.err()}"
assert first_result.ok_value, f"first channel_send returned an empty handle: {first_result.ok_value!r}"
second_result = node.channel_send(RC04_DISTINCT_CHANNEL_ID, create_message_bindings())
assert second_result.is_ok(), f"second channel_send failed: {second_result.err()}"
assert second_result.ok_value, f"second channel_send returned an empty handle: {second_result.ok_value!r}"
assert (
first_result.ok_value != second_result.ok_value
), f"channel_send handles must be distinct, got the same id twice: {first_result.ok_value!r}"
finally:
node.stop_and_destroy()

View File

@ -0,0 +1,93 @@
import base64
from src.steps.common import StepsCommon
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import (
EventCollector,
assert_event_invariants,
create_message_bindings,
get_node_multiaddr,
wait_for_connected,
wait_for_propagated,
wait_for_error,
)
# Payload above DefaultMaxWakuMessageSize (150KiB), so the relay publish
# rejects it instead of failing with NO_PEERS_TO_RELAY.
OVERSIZED_PAYLOAD_BYTES = 200 * 1024
ERROR_TIMEOUT_S = 30.0
MESSAGE_SIZE_EXCEEDED_MSG = "Message size exceeded"
class TestS13RelayHardFailureWithoutFallback(StepsCommon):
"""
S13: relay path is reachable (a relay peer is connected, so the publish
gets past NO_PEERS_TO_RELAY), but the relay publish fails for another
reason. An oversized payload is used so the relay processor rejects the
message immediately. No lightpush fallback is configured.
- Expected: Ok(RequestId), then a message_error event.
"""
def test_s13_relay_hard_failure_without_fallback(self, node_config):
sender_collector = EventCollector()
node_config.update(
{
"relay": True,
"numShardsInNetwork": 1,
}
)
sender_result = WrapperManager.create_and_start(
config=node_config,
event_cb=sender_collector.event_callback,
)
assert sender_result.is_ok(), f"Failed to start sender: {sender_result.err()}"
with sender_result.ok_value as sender_node:
relay_config = {
**node_config,
"staticnodes": [get_node_multiaddr(sender_node)],
"portsShift": 1,
}
relay_result = WrapperManager.create_and_start(config=relay_config)
assert relay_result.is_ok(), f"Failed to start relay peer: {relay_result.err()}"
with relay_result.ok_value:
# A connected relay peer means the publish gets past
# NO_PEERS_TO_RELAY and actually reaches the relay processor.
assert wait_for_connected(sender_collector) is not None, (
f"Sender did not reach Connected/PartiallyConnected. " f"Collected events: {sender_collector.events}"
)
oversized_payload = base64.b64encode(b"x" * OVERSIZED_PAYLOAD_BYTES).decode()
message = create_message_bindings(
payload=oversized_payload,
contentTopic="/test/1/s13-relay-hard-failure/proto",
)
send_result = sender_node.send_message(message=message)
assert send_result.is_ok(), f"send() must return Ok(RequestId), got: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
error_event = wait_for_error(
collector=sender_collector,
request_id=request_id,
timeout_s=ERROR_TIMEOUT_S,
)
assert error_event is not None, (
f"No message_error event within {ERROR_TIMEOUT_S}s from the " f"relay processor. Collected events: {sender_collector.events}"
)
assert error_event["requestId"] == request_id
assert MESSAGE_SIZE_EXCEEDED_MSG in (error_event.get("error") or ""), (
f"Expected error to contain {MESSAGE_SIZE_EXCEEDED_MSG!r}.\n" f"Got: {error_event.get('error')!r}\n" f"Full event: {error_event}"
)
propagated = wait_for_propagated(sender_collector, request_id, timeout_s=0)
assert propagated is None, f"Unexpected message_propagated event for a failed relay publish: {propagated}"
assert_event_invariants(sender_collector, request_id)

View File

@ -0,0 +1,326 @@
import json
import subprocess
import sys
import pytest
from src.steps.common import StepsCommon
from src.libs.common import to_base64
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import (
EventCollector,
assert_event_invariants,
create_message_bindings,
get_node_multiaddr,
wait_for_connected,
wait_for_propagated,
wait_for_sent,
wait_for_error,
)
PROPAGATED_TIMEOUT_S = 30.0
S01_EXPECTED_ERROR_FRAGMENT = "not initialized"
# Destroyed-handle path fails synchronously in the C layer (no callback),
# so the binding surfaces a different string than the nil-handle path.
S01_DESTROYED_HANDLE_ERROR_FRAGMENT = "immediate call failed"
S01_SUBPROCESS_TIMEOUT_S = 30
S01_RESULT_MARKER = "__S01_RESULT__"
SEND_AFTER_DESTROY_RESULT_MARKER = "__SEND_AFTER_DESTROY_RESULT__"
SEND_AFTER_DESTROY_SUBPROCESS_TIMEOUT_S = 60
S01_INVALID_HANDLE_HELPER = "tests.wrappers_tests.helpers.send_invalid_handle"
# S05: malformed content topics
S05_EXPECTED_ERROR_FRAGMENT = "Failed to auto-subscribe"
S05_MALFORMED_CONTENT_TOPICS = [
# No leading slash — parser rejects with "must start with slash".
("s05-invalid-no-leading-slash", "no-leading-slash"),
# Empty string — parser rejects empty content topic.
("", "empty"),
# Only 3 segments — content topics need /app/version/name/encoding.
("/app/1/name", "missing-encoding-segment"),
# Empty middle segment between slashes.
("/app//name/proto", "empty-middle-segment"),
]
class TestS01NilOrUninitializedHandle(StepsCommon):
"""S01 — send() on a nil/destroyed handle must Err, no events, no crash."""
@pytest.mark.skip(reason="see https://github.com/logos-messaging/logos-delivery/issues/3873")
def test_s01_send_on_uninitialized_handle(self):
completed = subprocess.run(
[sys.executable, "-m", S01_INVALID_HANDLE_HELPER, "nil", S01_RESULT_MARKER],
capture_output=True,
text=True,
timeout=S01_SUBPROCESS_TIMEOUT_S,
)
assert completed.returncode == 0, (
f"send() crashed on a nil handle (returncode={completed.returncode}). " f"stdout={completed.stdout!r} stderr={completed.stderr!r}"
)
result_line = next(
(l for l in completed.stdout.splitlines() if l.startswith(S01_RESULT_MARKER)),
None,
)
assert result_line, f"missing result marker. stdout={completed.stdout!r} stderr={completed.stderr!r}"
result = json.loads(result_line[len(S01_RESULT_MARKER) :])
assert result["is_ok"] is False, f"expected Err, got Ok({result['ok']!r})"
assert S01_EXPECTED_ERROR_FRAGMENT in (
result["err"] or ""
), f"expected error to mention {S01_EXPECTED_ERROR_FRAGMENT!r}, got: {result['err']!r}"
@pytest.mark.skip(reason="see https://github.com/logos-messaging/logos-delivery/issues/3863")
def test_s01_send_on_destroyed_handle(self):
completed = subprocess.run(
[sys.executable, "-m", S01_INVALID_HANDLE_HELPER, "destroyed", SEND_AFTER_DESTROY_RESULT_MARKER],
capture_output=True,
text=True,
timeout=SEND_AFTER_DESTROY_SUBPROCESS_TIMEOUT_S,
)
assert completed.returncode == 0, (
f"send() crashed on a destroyed handle (returncode={completed.returncode}). " f"stdout={completed.stdout!r} stderr={completed.stderr!r}"
)
result_line = next(
(l for l in completed.stdout.splitlines() if l.startswith(SEND_AFTER_DESTROY_RESULT_MARKER)),
None,
)
assert result_line, f"missing result marker. stdout={completed.stdout!r} stderr={completed.stderr!r}"
result = json.loads(result_line[len(SEND_AFTER_DESTROY_RESULT_MARKER) :])
assert result["stage"] == "send_message", f"setup failed at stage {result['stage']!r}: {result['err']!r}"
assert result["is_ok"] is False, f"expected Err, got Ok({result['ok']!r})"
err_msg = result["err"] or ""
assert S01_EXPECTED_ERROR_FRAGMENT in err_msg or S01_DESTROYED_HANDLE_ERROR_FRAGMENT in err_msg, (
f"expected error to mention {S01_EXPECTED_ERROR_FRAGMENT!r} " f"or {S01_DESTROYED_HANDLE_ERROR_FRAGMENT!r}, got: {result['err']!r}"
)
assert result["events_after_send"] == [], f"expected no events after send(), got: {result['events_after_send']}"
class TestS03SendOnAlreadySubscribedTopic(StepsCommon):
"""
S03 Send on already-subscribed content topic.
Sender explicitly calls subscribe_content_topic() before send().
The send path must behave identically to the auto-subscribe case:
Propagated arrives, no Sent (store disabled), no Error.
Topology mirrors S06 (relay-only sender + relay peer, no store).
Purpose: proves the send path is identical when auto-subscription is skipped.
"""
def test_s03_send_on_already_subscribed_content_topic(self, node_config):
sender_collector = EventCollector()
content_topic = "/test/1/s03-already-subscribed/proto"
node_config.update(
{
"relay": True,
"store": False,
"lightpush": False,
"filter": False,
"discv5Discovery": False,
"numShardsInNetwork": 1,
"reliabilityEnabled": True,
}
)
sender_result = WrapperManager.create_and_start(
config=node_config,
event_cb=sender_collector.event_callback,
)
assert sender_result.is_ok(), f"Failed to start sender: {sender_result.err()}"
with sender_result.ok_value as sender:
peer_config = {
**node_config,
"staticnodes": [get_node_multiaddr(sender)],
"portsShift": 1,
}
peer_result = WrapperManager.create_and_start(config=peer_config)
assert peer_result.is_ok(), f"Failed to start relay peer: {peer_result.err()}"
with peer_result.ok_value:
assert wait_for_connected(sender_collector) is not None, "Sender did not reach Connected/PartiallyConnected state"
# Explicit subscribe before send — this is what S03 is about.
# The send path must still return Ok(RequestId) and emit the
# same events as the auto-subscribe topology in S06.
subscribe_result = sender.subscribe_content_topic(content_topic)
assert subscribe_result.is_ok(), f"subscribe_content_topic failed: {subscribe_result.err()}"
message = create_message_bindings(
payload=to_base64("S03 already-subscribed test payload"),
contentTopic=content_topic,
)
send_result = sender.send_message(message=message)
assert send_result.is_ok(), f"send() failed: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
propagated = wait_for_propagated(
collector=sender_collector,
request_id=request_id,
timeout_s=PROPAGATED_TIMEOUT_S,
)
assert propagated is not None, (
f"No message_propagated event within {PROPAGATED_TIMEOUT_S}s. " f"Collected events: {sender_collector.events}"
)
assert propagated["requestId"] == request_id
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error event: {error}"
sent = wait_for_sent(sender_collector, request_id, timeout_s=0)
assert sent is None, f"Unexpected message_sent event (store is disabled): {sent}"
assert_event_invariants(sender_collector, request_id)
class TestS04UnsubscribeThenSendSameTopic(StepsCommon):
"""
S04 Unsubscribe, then send the same content topic again.
Sender subscribes to topic A, unsubscribes from A, then sends on A.
The send path must re-establish topic interest and deliver normally.
Topology mirrors S06 (relay-only sender + relay peer, no store).
Expected: send() returns Ok(RequestId), Propagated arrives,
no Sent (store disabled), no Error.
Purpose: verifies send() re-establishes topic interest after local unsubscribe.
"""
def test_s04_unsubscribe_then_send_same_content_topic(self, node_config):
sender_collector = EventCollector()
content_topic = "/test/1/s04-unsubscribe-resend/proto"
node_config.update(
{
"relay": True,
"store": False,
"lightpush": False,
"filter": False,
"discv5Discovery": False,
"numShardsInNetwork": 1,
"reliabilityEnabled": True,
}
)
sender_result = WrapperManager.create_and_start(
config=node_config,
event_cb=sender_collector.event_callback,
)
assert sender_result.is_ok(), f"Failed to start sender: {sender_result.err()}"
with sender_result.ok_value as sender:
peer_config = {
**node_config,
"staticnodes": [get_node_multiaddr(sender)],
"portsShift": 1,
}
peer_result = WrapperManager.create_and_start(config=peer_config)
assert peer_result.is_ok(), f"Failed to start relay peer: {peer_result.err()}"
with peer_result.ok_value:
assert wait_for_connected(sender_collector) is not None, "Sender did not reach Connected/PartiallyConnected state"
# subscribe -> unsubscribe -> send: send() must re-establish
# topic interest internally.
subscribe_result = sender.subscribe_content_topic(content_topic)
assert subscribe_result.is_ok(), f"subscribe_content_topic failed: {subscribe_result.err()}"
unsubscribe_result = sender.unsubscribe_content_topic(content_topic)
assert unsubscribe_result.is_ok(), f"unsubscribe_content_topic failed: {unsubscribe_result.err()}"
message = create_message_bindings(
payload=to_base64("S04 unsubscribe-then-send test payload"),
contentTopic=content_topic,
)
send_result = sender.send_message(message=message)
assert send_result.is_ok(), f"send() failed after unsubscribe: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
propagated = wait_for_propagated(
collector=sender_collector,
request_id=request_id,
timeout_s=PROPAGATED_TIMEOUT_S,
)
assert propagated is not None, (
f"No message_propagated event within {PROPAGATED_TIMEOUT_S}s "
f"after unsubscribe + send. Collected events: {sender_collector.events}"
)
assert propagated["requestId"] == request_id
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error event: {error}"
sent = wait_for_sent(sender_collector, request_id, timeout_s=0)
assert sent is None, f"Unexpected message_sent event (store is disabled): {sent}"
assert_event_invariants(sender_collector, request_id)
class TestS05AutoSubscribeFailureBeforeTaskCreation(StepsCommon):
"""
S05 Auto-subscribe failure before task creation.
Sender is initialized but auto-subscription is forced to fail by using a
malformed content topic that breaks shard resolution inside
SubscriptionManager.subscribe().
Expected: send() returns Err with an "auto-subscribe" message, no events.
Purpose: covers the last synchronous error path before request ID creation,
across several distinct validator branches.
"""
@pytest.mark.parametrize(
"content_topic",
[topic for topic, _ in S05_MALFORMED_CONTENT_TOPICS],
ids=[case_id for _, case_id in S05_MALFORMED_CONTENT_TOPICS],
)
def test_s05_send_fails_when_auto_subscribe_fails(self, node_config, content_topic):
sender_collector = EventCollector()
node_config.update(
{
"relay": True,
"numShardsInNetwork": 1,
}
)
sender_result = WrapperManager.create_and_start(
config=node_config,
event_cb=sender_collector.event_callback,
)
assert sender_result.is_ok(), f"Failed to start sender: {sender_result.err()}"
with sender_result.ok_value as sender:
# Malformed content topic — SubscriptionManager.subscribe() cannot
# resolve it to a shard, so auto-subscribe inside send() fails.
message = create_message_bindings(
payload=to_base64("S05 auto-subscribe failure payload"),
contentTopic=content_topic,
)
send_result = sender.send_message(message=message)
assert send_result.is_err(), (
f"send() must return Err when auto-subscribe fails for " f"content_topic={content_topic!r}, got Ok({send_result.ok_value!r})"
)
error_message = send_result.err() or ""
assert S05_EXPECTED_ERROR_FRAGMENT in error_message, (
f"expected error to mention {S05_EXPECTED_ERROR_FRAGMENT!r} " f"for content_topic={content_topic!r}, got: {error_message!r}"
)
# No request id was created, so no events should be emitted.
assert sender_collector.events == [] or all(
event.get("eventType") == "connection_status_change" for event in sender_collector.events
), f"Unexpected events after a pre-task-creation failure: {sender_collector.events}"

View File

@ -0,0 +1,500 @@
import pytest
from src.env_vars import NODE_1
from src.steps.common import StepsCommon
from src.libs.common import delay, to_base64
from src.node.waku_node import WakuNode
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import (
EventCollector,
assert_event_invariants,
create_message_bindings,
get_node_multiaddr,
wait_for_propagated,
wait_for_sent,
wait_for_error,
)
from src.test_data import DEFAULT_CLUSTER_ID
from tests.wrappers_tests.conftest import build_node_config
PROPAGATED_TIMEOUT_S = 30.0
NO_SENT_OBSERVATION_S = 5.0
SENT_AFTER_STORE_TIMEOUT_S = 60.0
RECOVERY_TIMEOUT_S = 45.0
SERVICE_DOWN_SETTLE_S = 3.0
# Default-cluster shard-0 pubsub topic; used to subscribe the S11 docker store
# peer so it joins the same relay mesh as the wrapper nodes (wrapper config
# uses numShardsInNetwork=1 => shard 0).
STORE_PEER_PUBSUB_TOPIC = f"/waku/2/rs/{DEFAULT_CLUSTER_ID}/0"
class TestS11EdgeSenderLightpushAndStore(StepsCommon):
"""
S11 Edge sender with lightpush path and store validation.
Edge sender has no local relay; it publishes via a wrapper lightpush peer
and validates delivery via a docker store peer. Reliability enabled.
Topology:
[LightpushPeer] wrapper, relay=True, lightpush=True, store=False
[StorePeer] docker WakuNode, relay=true, store=true,
dials the lightpush peer via add_peers and subscribes
to the same shard-0 pubsub topic so it joins the
relay mesh and archives propagated messages.
[Edge] wrapper, mode="Edge",
staticnodes=[lightpush_peer],
storenode=store_peer,
reliabilityEnabled=True
Expected: send() returns Ok(RequestId), Propagated arrives, then Sent
(store validation succeeds), no Error.
Purpose: edge-mode fully validated success path against a real docker
store node (cross-implementation check).
"""
@pytest.mark.docker_required
def test_s11_edge_lightpush_with_store_validation(self):
sender_collector = EventCollector()
common = {
"filter": False,
"discv5Discovery": True,
"numShardsInNetwork": 1,
}
lightpush_config = build_node_config(
relay=True,
lightpush=True,
store=False,
**common,
)
lightpush_result = WrapperManager.create_and_start(config=lightpush_config)
assert lightpush_result.is_ok(), f"Failed to start lightpush peer: {lightpush_result.err()}"
with lightpush_result.ok_value as lightpush_peer:
lightpush_multiaddr = get_node_multiaddr(lightpush_peer)
# Docker store peer — real nwaku node running as the store backend.
# Dial the lightpush peer via add_peers and subscribe to the same
# shard-0 pubsub topic so it joins the relay mesh and archives
# messages propagated by the lightpush peer.
store_peer = WakuNode(NODE_1, f"s11_store_{self.test_id}")
store_peer.start(relay="true", store="true")
self.add_node_peer(store_peer, [lightpush_multiaddr])
store_peer.set_relay_subscriptions([STORE_PEER_PUBSUB_TOPIC])
store_multiaddr = store_peer.get_multiaddr_with_id()
edge_config = build_node_config(
mode="Edge",
relay=False,
store=False,
# assuming disc5v already happened
staticnodes=[lightpush_multiaddr, store_multiaddr],
reliabilityEnabled=True,
**common,
)
edge_result = WrapperManager.create_and_start(
config=edge_config,
event_cb=sender_collector.event_callback,
)
assert edge_result.is_ok(), f"Failed to start edge sender: {edge_result.err()}"
with edge_result.ok_value as edge_sender:
message = create_message_bindings(
payload=to_base64("S11 edge lightpush + store test payload"),
contentTopic="/test/1/s11-edge-lightpush-store/proto",
)
send_result = edge_sender.send_message(message=message)
assert send_result.is_ok(), f"send() failed: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
propagated = wait_for_propagated(
collector=sender_collector,
request_id=request_id,
timeout_s=PROPAGATED_TIMEOUT_S,
)
assert propagated is not None, (
f"No message_propagated event within {PROPAGATED_TIMEOUT_S}s. " f"Collected events: {sender_collector.events}"
)
assert propagated["requestId"] == request_id
sent = wait_for_sent(
collector=sender_collector,
request_id=request_id,
timeout_s=SENT_AFTER_STORE_TIMEOUT_S,
)
assert sent is not None, (
f"No message_sent event within {SENT_AFTER_STORE_TIMEOUT_S}s " f"after propagation. Collected events: {sender_collector.events}"
)
assert sent["requestId"] == request_id
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error event: {error}"
assert_event_invariants(sender_collector, request_id)
class TestS16LightpushPeerAppearsLater(StepsCommon):
"""
S16 No delivery peers at T0, lightpush peer appears later.
The edge sender has the lightpush service in its staticnodes, but the
service is stopped before the sender starts, so there is no reachable
delivery peer at T0. send() is called while the service is down. The
service is restarted during the retry window; the sender connects to it
and a later retry delivers the message.
Expected: send() returns Ok(RequestId), then eventually Propagated.
"""
@pytest.mark.xfail(reason="binding cannot restart a node or add peers at runtime")
def test_s16_lightpush_peer_appears_later(self):
sender_collector = EventCollector()
common = {
"store": False,
"filter": False,
"discv5Discovery": False,
"numShardsInNetwork": 1,
}
# Start the lightpush service once to obtain its multiaddr, then stop
# it so the sender has no reachable peer at T0. The same node object
# is restarted later, so the address stays valid.
service_config = build_node_config(relay=True, lightpush=True, **common)
service_result = WrapperManager.create_and_start(config=service_config)
assert service_result.is_ok(), f"Failed to start lightpush peer: {service_result.err()}"
with service_result.ok_value as service:
service_multiaddr = get_node_multiaddr(service)
stop_result = service.stop_node()
assert stop_result.is_ok(), f"Failed to stop lightpush peer: {stop_result.err()}"
delay(SERVICE_DOWN_SETTLE_S)
# Edge sender is a lightpush client; its only peer is the service,
# which is currently down.
edge_config = build_node_config(
mode="Edge",
relay=False,
staticnodes=[service_multiaddr],
**common,
)
edge_result = WrapperManager.create_and_start(
config=edge_config,
event_cb=sender_collector.event_callback,
)
assert edge_result.is_ok(), f"Failed to start edge sender: {edge_result.err()}"
with edge_result.ok_value as edge_sender:
# send() is invoked while the service is down.
msg = create_message_bindings(
payload=to_base64("S16 lightpush peer appears later"),
contentTopic="/test/1/s16-late-lightpush/proto",
)
send_result = edge_sender.send_message(message=msg)
assert send_result.is_ok(), f"send() failed: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
delay(SERVICE_DOWN_SETTLE_S)
early_propagated = wait_for_propagated(sender_collector, request_id, timeout_s=0)
assert early_propagated is None, f"message_propagated arrived before the lightpush peer was reachable: {early_propagated}"
# The lightpush peer comes back during the retry window.
restart_result = service.start_node()
assert restart_result.is_ok(), f"Failed to restart lightpush peer: {restart_result.err()}"
propagated = wait_for_propagated(
collector=sender_collector,
request_id=request_id,
timeout_s=RECOVERY_TIMEOUT_S,
)
assert propagated is not None, (
f"No message_propagated within {RECOVERY_TIMEOUT_S}s "
f"after the lightpush peer joined. "
f"Collected events: {sender_collector.events}"
)
assert propagated["requestId"] == request_id
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error after recovery: {error}"
assert_event_invariants(sender_collector, request_id)
class TestS18StagedTopologyReformation(StepsCommon):
"""
S18: relay absent at T0, lightpush absent at T0, both appear in stages.
The sender has relay and lightpush enabled but starts fully isolated:
no relay peers and no lightpush peers. send() is called before any
peer exists. Peers then appear one stage at a time. The message must
succeed as soon as any valid path becomes available.
Two orderings are covered, one per test method:
- lightpush appears first, then relay.
- relay appears first, then lightpush.
Both prove the send service recovers from arbitrary staged topology
reformation.
Common topology per stage:
[Sender] relay=True, lightpush=True, isolated at T0.
[LightpushPeer] relay=True, lightpush=True, staticnodes=[sender].
[RelayPeer] relay=True, staticnodes=[sender].
"""
# Common config shared by the sender and both staged peers.
_COMMON = {
"relay": True,
"lightpush": True,
"store": False,
"filter": False,
"discv5Discovery": False,
"numShardsInNetwork": 1,
}
def test_s18_lightpush_first_then_relay(self, node_config):
"""Sender isolated at T0; lightpush peer appears, then relay peer."""
sender_collector = EventCollector()
node_config.update(self._COMMON)
sender_result = WrapperManager.create_and_start(
config=node_config,
event_cb=sender_collector.event_callback,
)
assert sender_result.is_ok(), f"Failed to start sender: {sender_result.err()}"
with sender_result.ok_value as sender:
# send() before any peer exists must still return Ok(RequestId).
message = create_message_bindings(
payload=to_base64("S18 lightpush-first staged recovery"),
contentTopic="/test/1/s18-lightpush-first/proto",
)
send_result = sender.send_message(message=message)
assert send_result.is_ok(), f"send() must return Ok(RequestId) while isolated, got: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
# No peers yet: the message must not propagate.
delay(SERVICE_DOWN_SETTLE_S)
early_propagated = wait_for_propagated(sender_collector, request_id, timeout_s=0)
assert early_propagated is None, f"message_propagated arrived before any peer joined: {early_propagated}"
# Stage 1: lightpush peer appears.
lightpush_config = {
**node_config,
"staticnodes": [get_node_multiaddr(sender)],
"portsShift": 1,
}
lightpush_result = WrapperManager.create_and_start(config=lightpush_config)
assert lightpush_result.is_ok(), f"Failed to start lightpush peer: {lightpush_result.err()}"
with lightpush_result.ok_value:
# Stage 2: relay peer appears.
relay_config = {
**node_config,
"lightpush": False,
"staticnodes": [get_node_multiaddr(sender)],
"portsShift": 2,
}
relay_result = WrapperManager.create_and_start(config=relay_config)
assert relay_result.is_ok(), f"Failed to start relay peer: {relay_result.err()}"
with relay_result.ok_value:
# The message must succeed once any valid path is available.
propagated = wait_for_propagated(
collector=sender_collector,
request_id=request_id,
timeout_s=RECOVERY_TIMEOUT_S,
)
assert propagated is not None, (
f"No message_propagated within {RECOVERY_TIMEOUT_S}s "
f"after peers appeared in stages. "
f"Collected events: {sender_collector.events}"
)
assert propagated["requestId"] == request_id
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error after staged recovery: {error}"
assert_event_invariants(sender_collector, request_id)
def test_s18_relay_first_then_lightpush(self, node_config):
"""Sender isolated at T0; relay peer appears, then lightpush peer."""
sender_collector = EventCollector()
node_config.update(self._COMMON)
sender_result = WrapperManager.create_and_start(
config=node_config,
event_cb=sender_collector.event_callback,
)
assert sender_result.is_ok(), f"Failed to start sender: {sender_result.err()}"
with sender_result.ok_value as sender:
# send() before any peer exists must still return Ok(RequestId).
message = create_message_bindings(
payload=to_base64("S18 relay-first staged recovery"),
contentTopic="/test/1/s18-relay-first/proto",
)
send_result = sender.send_message(message=message)
assert send_result.is_ok(), f"send() must return Ok(RequestId) while isolated, got: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
# No peers yet: the message must not propagate.
delay(SERVICE_DOWN_SETTLE_S)
early_propagated = wait_for_propagated(sender_collector, request_id, timeout_s=0)
assert early_propagated is None, f"message_propagated arrived before any peer joined: {early_propagated}"
# Stage 1: relay peer appears.
relay_config = {
**node_config,
"lightpush": False,
"staticnodes": [get_node_multiaddr(sender)],
"portsShift": 1,
}
relay_result = WrapperManager.create_and_start(config=relay_config)
assert relay_result.is_ok(), f"Failed to start relay peer: {relay_result.err()}"
with relay_result.ok_value:
# Stage 2: lightpush peer appears.
lightpush_config = {
**node_config,
"staticnodes": [get_node_multiaddr(sender)],
"portsShift": 2,
}
lightpush_result = WrapperManager.create_and_start(config=lightpush_config)
assert lightpush_result.is_ok(), f"Failed to start lightpush peer: {lightpush_result.err()}"
with lightpush_result.ok_value:
# The message must succeed once any valid path is available.
propagated = wait_for_propagated(
collector=sender_collector,
request_id=request_id,
timeout_s=RECOVERY_TIMEOUT_S,
)
assert propagated is not None, (
f"No message_propagated within {RECOVERY_TIMEOUT_S}s "
f"after peers appeared in stages. "
f"Collected events: {sender_collector.events}"
)
assert propagated["requestId"] == request_id
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error after staged recovery: {error}"
assert_event_invariants(sender_collector, request_id)
class TestS25EphemeralLightpushWithStore(StepsCommon):
"""
S25 Ephemeral message over lightpush with a reachable store peer.
The sender is an Edge node, so it has no local relay and publishes only
via lightpush. A docker store peer is reachable and joined to the
lightpush peer's relay mesh, so the message does reach a store.
Because ephemeral messages are never store-validated, the expected result
is Propagated only no Sent even though the store peer is reachable.
This is the lightpush-transport counterpart of S24 and proves the
ephemeral rule is transport-independent.
Topology mirrors S11:
[LightpushPeer] wrapper, relay=True, lightpush=True, store=False
[StorePeer] docker WakuNode, relay=true, store=true, joined to
the lightpush peer's shard-0 relay mesh
[Edge] wrapper, mode="Edge", reliabilityEnabled=True,
staticnodes=[lightpush_peer, store_peer]
"""
@pytest.mark.docker_required
def test_s25_ephemeral_lightpush_with_store(self):
sender_collector = EventCollector()
common = {
"store": False,
"numShardsInNetwork": 1,
}
lightpush_config = build_node_config(
relay=True,
lightpush=True,
**common,
)
lightpush_result = WrapperManager.create_and_start(config=lightpush_config)
assert lightpush_result.is_ok(), f"Failed to start lightpush peer: {lightpush_result.err()}"
with lightpush_result.ok_value as lightpush_peer:
lightpush_multiaddr = get_node_multiaddr(lightpush_peer)
# Docker store peer joined to the lightpush peer's shard-0 relay
# mesh, so messages propagated by the lightpush peer are archived.
store_peer = WakuNode(NODE_1, f"s25_store_{self.test_id}")
store_peer.start(relay="true", store="true")
self.add_node_peer(store_peer, [lightpush_multiaddr])
store_peer.set_relay_subscriptions([STORE_PEER_PUBSUB_TOPIC])
store_multiaddr = store_peer.get_multiaddr_with_id()
edge_config = build_node_config(
mode="Edge",
relay=False,
staticnodes=[lightpush_multiaddr, store_multiaddr],
**common,
)
edge_result = WrapperManager.create_and_start(
config=edge_config,
event_cb=sender_collector.event_callback,
)
assert edge_result.is_ok(), f"Failed to start edge sender: {edge_result.err()}"
with edge_result.ok_value as edge_sender:
message = create_message_bindings(
payload=to_base64("S25 ephemeral lightpush + store payload"),
contentTopic="/test/1/s25-ephemeral-lightpush/proto",
ephemeral=True,
)
send_result = edge_sender.send_message(message=message)
assert send_result.is_ok(), f"send() failed: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send() returned an empty RequestId"
propagated = wait_for_propagated(
collector=sender_collector,
request_id=request_id,
timeout_s=PROPAGATED_TIMEOUT_S,
)
assert propagated is not None, (
f"No message_propagated event within {PROPAGATED_TIMEOUT_S}s. " f"Collected events: {sender_collector.events}"
)
assert propagated["requestId"] == request_id
# Ephemeral messages are never store-validated, so no Sent
# event must arrive even though the store peer is reachable.
sent = wait_for_sent(
collector=sender_collector,
request_id=request_id,
timeout_s=NO_SENT_OBSERVATION_S,
)
assert sent is None, (
f"Unexpected message_sent event for an ephemeral message. "
f"Ephemeral messages must never be store-validated.\n"
f"Sent event: {sent}\n"
f"Collected events: {sender_collector.events}"
)
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error event: {error}"
assert_event_invariants(sender_collector, request_id)

View File

@ -0,0 +1,117 @@
from src.steps.common import StepsCommon
from src.libs.common import to_base64
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import (
EventCollector,
assert_event_invariants,
create_message_bindings,
get_node_multiaddr,
wait_for_connected,
wait_for_propagated,
wait_for_error,
)
from src.test_data import CONTENT_TOPICS_DIFFERENT_SHARDS
PROPAGATED_TIMEOUT_S = 30.0
class TestS29SendOnTopicsMappingToDifferentShards(StepsCommon):
"""
S29 Send on two different content topics that map to different shards.
Sender has a relay peer reachable on shard X and shard Y; topic A maps to
shard X and topic B maps to shard Y. Two independent sends, one per topic.
Expected: both sends return Ok(RequestId), and each request gets its own
message_propagated event following the availability of its own shard.
Purpose: ensures shard derivation and delivery behavior are topic-specific.
"""
# Topic A -> shard 0, Topic B -> shard 1 (per CONTENT_TOPICS_DIFFERENT_SHARDS).
TOPIC_A = CONTENT_TOPICS_DIFFERENT_SHARDS[0]
TOPIC_B = CONTENT_TOPICS_DIFFERENT_SHARDS[1]
def test_s29_send_on_topics_mapping_to_different_shards(self, node_config):
sender_collector = EventCollector()
# numShardsInNetwork=8 so the two topics resolve to distinct shards
# (shard 0 and shard 1) instead of being collapsed onto shard 0.
node_config.update(
{
"relay": True,
"store": False,
"lightpush": False,
"filter": False,
"discv5Discovery": False,
"numShardsInNetwork": 8,
"reliabilityEnabled": True,
}
)
sender_result = WrapperManager.create_and_start(
config=node_config,
event_cb=sender_collector.event_callback,
)
assert sender_result.is_ok(), f"Failed to start sender: {sender_result.err()}"
with sender_result.ok_value as sender:
peer_config = {
**node_config,
"staticnodes": [get_node_multiaddr(sender)],
"portsShift": 1,
}
peer_result = WrapperManager.create_and_start(config=peer_config)
assert peer_result.is_ok(), f"Failed to start relay peer: {peer_result.err()}"
with peer_result.ok_value:
assert wait_for_connected(sender_collector) is not None, "Sender did not reach Connected/PartiallyConnected state"
message_a = create_message_bindings(
payload=to_base64("S29 shard X payload"),
contentTopic=self.TOPIC_A,
)
send_a = sender.send_message(message=message_a)
assert send_a.is_ok(), f"send() on TOPIC_A failed: {send_a.err()}"
request_id_a = send_a.ok_value
assert request_id_a, "send() on TOPIC_A returned an empty RequestId"
# Send on topic B (shard Y).
message_b = create_message_bindings(
payload=to_base64("S29 shard Y payload"),
contentTopic=self.TOPIC_B,
)
send_b = sender.send_message(message=message_b)
assert send_b.is_ok(), f"send() on TOPIC_B failed: {send_b.err()}"
request_id_b = send_b.ok_value
assert request_id_b, "send() on TOPIC_B returned an empty RequestId"
assert request_id_a != request_id_b, "Each send must produce a distinct RequestId"
# Each request propagates over its own shard's mesh independently.
propagated_a = wait_for_propagated(
collector=sender_collector,
request_id=request_id_a,
timeout_s=PROPAGATED_TIMEOUT_S,
)
assert propagated_a is not None, (
f"No message_propagated event for TOPIC_A within {PROPAGATED_TIMEOUT_S}s. " f"Collected events: {sender_collector.events}"
)
assert propagated_a["requestId"] == request_id_a
propagated_b = wait_for_propagated(
collector=sender_collector,
request_id=request_id_b,
timeout_s=PROPAGATED_TIMEOUT_S,
)
assert propagated_b is not None, (
f"No message_propagated event for TOPIC_B within {PROPAGATED_TIMEOUT_S}s. " f"Collected events: {sender_collector.events}"
)
assert propagated_b["requestId"] == request_id_b
# No cross-talk: neither request should produce an error.
for request_id in (request_id_a, request_id_b):
error = wait_for_error(sender_collector, request_id, timeout_s=0)
assert error is None, f"Unexpected message_error event for {request_id}: {error}"
assert_event_invariants(sender_collector, request_id_a)
assert_event_invariants(sender_collector, request_id_b)

View File

@ -0,0 +1,130 @@
import pytest
from src.steps.common import StepsCommon
from src.libs.custom_logger import get_custom_logger
from src.node.wrappers_manager import WrapperManager
from src.node.wrapper_helpers import (
EventCollector,
create_message_bindings,
enr_udp_port,
get_node_bound_ports,
get_node_multiaddr,
get_node_tcp_port,
wait_for_propagated,
)
from tests.wrappers_tests.conftest import build_node_config, free_port
logger = get_custom_logger(__name__)
PROPAGATED_TIMEOUT_S = 30.0
# The five service ports covered by logos-messaging/logos-delivery#3828:
# (config field, key in MyBoundPorts, config to enable the service, default port)
SERVICE_PORTS = [
("tcpPort", "tcp", {}, 60000),
("discv5UdpPort", "discv5Udp", {"discv5Discovery": True}, 9000),
("websocketPort", "webSocket", {"websocketSupport": True}, 8000),
("restPort", "rest", {"rest": True}, 8645),
("metricsServerPort", "metrics", {"metricsServer": True}, 8008),
]
class TestWrapperAutoPortAllocation(StepsCommon):
"""Corner case: port 0 triggers auto-port allocation.
Tracks logos-messaging/logos-delivery#3828:
- any service port set to 0 gets a free port auto-assigned at bind time
- defaults remain concrete, so auto-port is opt-in via an explicit 0
- bound values are exposed through MyBoundPorts (0 = service disabled)
- the ENR is rebuilt after discv5 startup, so it advertises the
actually-bound UDP port instead of the configured 0
"""
def test_auto_port_starts_node_with_tcp_and_discv5_zero(self):
config = build_node_config(tcpPort=0, discv5UdpPort=0)
result = WrapperManager.create_and_start(config=config)
assert result.is_ok(), f"create_and_start failed: {result.err()}"
with result.ok_value as node:
tcp_port = get_node_tcp_port(node)
assert tcp_port != 0, "multiaddr still reports port 0; auto-port did not happen"
assert get_node_bound_ports(node)["tcp"] == tcp_port, "MyBoundPorts disagrees with the multiaddr port"
def test_auto_port_node_can_propagate_message(self):
# End-to-end: two auto-port nodes exchange a message.
# numShardsInNetwork=1 enables autosharding, required by the send API.
collector = EventCollector()
sender_config = build_node_config(tcpPort=0, discv5UdpPort=0, numShardsInNetwork=1)
sender_result = WrapperManager.create_and_start(config=sender_config, event_cb=collector.event_callback)
assert sender_result.is_ok(), f"sender start failed: {sender_result.err()}"
with sender_result.ok_value as sender:
peer_config = build_node_config(
tcpPort=0,
discv5UdpPort=0,
numShardsInNetwork=1,
staticnodes=[get_node_multiaddr(sender)],
)
peer_result = WrapperManager.create_and_start(config=peer_config)
assert peer_result.is_ok(), f"peer start failed: {peer_result.err()}"
with peer_result.ok_value:
send_result = sender.send_message(message=create_message_bindings())
assert send_result.is_ok(), f"send failed: {send_result.err()}"
request_id = send_result.ok_value
assert request_id, "send returned empty RequestId"
propagated = wait_for_propagated(collector, request_id, PROPAGATED_TIMEOUT_S)
assert propagated is not None, f"no message_propagated event. Events: {collector.events}"
@pytest.mark.parametrize("port_field, bound_key, enable_service, default_port", SERVICE_PORTS)
def test_auto_port_per_field(self, port_field, bound_key, enable_service, default_port):
# Each service port set to 0 in isolation, with its service enabled.
config = build_node_config(**{port_field: 0}, **enable_service)
result = WrapperManager.create_and_start(config=config)
assert result.is_ok(), f"create_and_start failed with {port_field}=0: {result.err()}"
with result.ok_value as node:
port = get_node_bound_ports(node)[bound_key]
assert port != 0, f"{port_field}=0 but nothing bound; auto-port did not happen"
assert port != default_port, f"{port_field}=0 bound the default {default_port}; expected an ephemeral port"
def test_concrete_ports_are_respected(self):
# Auto-port is opt-in: non-zero ports must bind exactly as requested.
requested = {field: free_port() for field, _, _, _ in SERVICE_PORTS}
enable_all = {key: value for _, _, enable, _ in SERVICE_PORTS for key, value in enable.items()}
config = build_node_config(**requested, **enable_all)
result = WrapperManager.create_and_start(config=config)
assert result.is_ok(), f"create_and_start failed: {result.err()}"
with result.ok_value as node:
bound = get_node_bound_ports(node)
for field, bound_key, _, _ in SERVICE_PORTS:
assert bound[bound_key] == requested[field], f"{field}: requested {requested[field]}, bound {bound[bound_key]}"
def test_bound_ports_zero_for_disabled_services(self):
# build_node_config leaves websocket, REST, metrics and discv5 off.
result = WrapperManager.create_and_start(config=build_node_config())
assert result.is_ok(), f"create_and_start failed: {result.err()}"
with result.ok_value as node:
bound = get_node_bound_ports(node)
assert bound["tcp"] != 0, "tcp is enabled but reports port 0"
for key in ("webSocket", "rest", "discv5Udp", "metrics"):
assert bound[key] == 0, f"{key} is disabled but reports port {bound[key]}"
def test_enr_advertises_bound_discv5_port(self):
# #3828 rebuilds the ENR after discv5 startup so it advertises the
# actually-bound UDP port, not the configured 0.
config = build_node_config(discv5UdpPort=0, discv5Discovery=True)
result = WrapperManager.create_and_start(config=config)
assert result.is_ok(), f"create_and_start failed: {result.err()}"
with result.ok_value as node:
discv5_port = get_node_bound_ports(node)["discv5Udp"]
assert discv5_port != 0, "discv5 enabled with port 0 but nothing bound"
enr_result = node.get_node_info("MyENR")
assert enr_result.is_ok(), f"MyENR query failed: {enr_result.err()}"
assert enr_udp_port(enr_result.ok_value.strip()) == discv5_port, "ENR was not rebuilt after discv5 startup"

View File

@ -55,6 +55,18 @@ int logosdelivery_remove_event_listener(
void *ctx,
uint64_t listenerId
);
typedef struct {
const char *channelIdStr;
const char *contentTopicStr;
const char *senderIdStr;
} ChannelCreateReq;
typedef struct { const char *channelIdStr; const char *messageJson; } ChannelSendReq;
typedef struct { const char *channelIdStr; } ChannelCloseReq;
int logosdelivery_channel_create(void *ctx, ReplyFn onReply, void *userData, const ChannelCreateReq *req);
int logosdelivery_channel_send(void *ctx, ReplyFn onReply, void *userData, const ChannelSendReq *req);
int logosdelivery_channel_close(void *ctx, ReplyFn onReply, void *userData, const ChannelCloseReq *req);
"""
)
@ -409,3 +421,94 @@ class NodeWrapper:
return Ok(result)
def destroy_keep_ctx(self, *, timeout_s: float = 20.0) -> Result[int, str]:
"""Destroy the node without nilling self.ctx afterwards.
Lets a library-contract test reach the C side with a dangling-but-non-nil
pointer, instead of relying on the binding's defensive nil-out.
"""
rc = lib.logosdelivery_destroy(self.ctx)
if rc != 0:
return Err(f"destroy_keep_ctx: call failed (ret={rc})")
return Ok(rc)
def channel_create(
self,
channel_id: str,
content_topic: str,
sender_id: str,
*,
timeout_s: float = 20.0,
) -> Result[str, str]:
state = _new_cb_state()
cb = self._make_waiting_reply_cb(state)
channel_buffer = ffi.new("char[]", channel_id.encode("utf-8"))
topic_buffer = ffi.new("char[]", content_topic.encode("utf-8"))
sender_buffer = ffi.new("char[]", sender_id.encode("utf-8"))
req = ffi.new(
"ChannelCreateReq *",
{
"channelIdStr": channel_buffer,
"contentTopicStr": topic_buffer,
"senderIdStr": sender_buffer,
},
)
rc = lib.logosdelivery_channel_create(self.ctx, cb, ffi.NULL, req)
if rc != 0:
return Err(f"channel_create: immediate call failed (ret={rc})")
wait_result = _wait_cb_raw(state, f"channel_create({channel_id})", timeout_s)
if wait_result.is_err():
return Err(wait_result.err())
cb_ret, cb_msg = wait_result.ok_value
if cb_ret != 0:
return Err(cb_msg.decode("utf-8") if cb_msg else f"channel_create({channel_id}): callback failed (ret={cb_ret})")
return Ok(cb_msg.decode("utf-8") if cb_msg else "")
def channel_send(self, channel_id: str, message: dict, *, timeout_s: float = 20.0) -> Result[str, str]:
state = _new_cb_state()
cb = self._make_waiting_reply_cb(state)
message_json = json.dumps(message, separators=(",", ":"), ensure_ascii=False)
channel_buffer = ffi.new("char[]", channel_id.encode("utf-8"))
message_buffer = ffi.new("char[]", message_json.encode("utf-8"))
req = ffi.new("ChannelSendReq *", {"channelIdStr": channel_buffer, "messageJson": message_buffer})
rc = lib.logosdelivery_channel_send(self.ctx, cb, ffi.NULL, req)
if rc != 0:
return Err(f"channel_send: immediate call failed (ret={rc})")
wait_result = _wait_cb_raw(state, f"channel_send({channel_id})", timeout_s)
if wait_result.is_err():
return Err(wait_result.err())
cb_ret, cb_msg = wait_result.ok_value
if cb_ret != 0:
return Err(cb_msg.decode("utf-8") if cb_msg else f"channel_send({channel_id}): callback failed (ret={cb_ret})")
return Ok(cb_msg.decode("utf-8") if cb_msg else "")
def channel_close(self, channel_id: str, *, timeout_s: float = 20.0) -> Result[str, str]:
state = _new_cb_state()
cb = self._make_waiting_reply_cb(state)
channel_buffer = ffi.new("char[]", channel_id.encode("utf-8"))
req = ffi.new("ChannelCloseReq *", {"channelIdStr": channel_buffer})
rc = lib.logosdelivery_channel_close(self.ctx, cb, ffi.NULL, req)
if rc != 0:
return Err(f"channel_close: immediate call failed (ret={rc})")
wait_result = _wait_cb_raw(state, f"channel_close({channel_id})", timeout_s)
if wait_result.is_err():
return Err(wait_result.err())
cb_ret, cb_msg = wait_result.ok_value
if cb_ret != 0:
return Err(cb_msg.decode("utf-8") if cb_msg else f"channel_close({channel_id}): callback failed (ret={cb_ret})")
return Ok(cb_msg.decode("utf-8") if cb_msg else "")