test(e2e): port the remaining channel delivery scenarios

Six channel tests were left behind in the interop repo when the suite
moved here, so the paths they cover — ingress filtering by content topic,
causal ordering, late joiners, three-way fan-out and channel re-creation —
are unguarded in this repo's CI.

Bring them over: RC04 send-after-close, RC07 foreign content topic, RC08
bidirectional causal order, RC09 late-joining receiver, RC12 three
participants, RC13 close/re-create then send.

RC07 and RC08 get store off like every other channel test; the interop
copies left the fixture default on, which pins an archive to the shared
working directory for no reason of theirs.
This commit is contained in:
Egor Rachkovskii
2026-08-17 15:26:25 +01:00
parent 4b4e2996a3
commit f2c64a3a8e
2 changed files with 394 additions and 0 deletions
@@ -23,6 +23,26 @@ MESSAGE_RECEIVED_EVENT = "message_received"
RC06_CHANNEL_PREFIX = "rc06-channel"
RC06_CONTENT_TOPIC = "/test/1/rc06-channel/proto"
RC07_CHANNEL_PREFIX = "rc07-channel"
RC07_CONTENT_TOPIC = "/test/1/rc07-channel/proto"
RC07_OTHER_CONTENT_TOPIC = "/test/1/rc07-other/proto"
RC08_CHANNEL_PREFIX = "rc08-channel"
RC08_CONTENT_TOPIC = "/test/1/rc08-channel/proto"
# A reply references what it answers only once that landed in the sender's SDS
# history; replying immediately races that write and asserts nothing.
RC08_CAUSAL_SETTLE_S = 5
RC09_CHANNEL_PREFIX = "rc09-channel"
RC09_CONTENT_TOPIC = "/test/1/rc09-channel/proto"
RC12_CHANNEL_PREFIX = "rc12-channel"
RC12_CONTENT_TOPIC = "/test/1/rc12-channel/proto"
SENDER_C = "rc12-sender-c"
RC13_CHANNEL_PREFIX = "rc13-channel"
RC13_CONTENT_TOPIC = "/test/1/rc13-channel/proto"
CLOSED_CHANNEL_PREFIX = "rc-closed-channel"
CLOSED_CONTENT_TOPIC = "/test/1/rc-closed-channel/proto"
@@ -33,6 +53,10 @@ DELIVERY_TIMEOUT_S = 50.0
# just a short grace window to catch any late channel_message_received.
NO_CHANNEL_DELIVERY_WINDOW_S = 10.0
# A has no peer to dial before it sends, so no mesh to settle.
RC09_SENDER_SETTLE_S = 2
RC09_RECOVERY_TIMEOUT_S = 90.0
def wait_for_channel_received(collector, channel_id, timeout_s, poll_interval_s=0.5):
deadline = time.monotonic() + timeout_s
@@ -45,6 +69,25 @@ def wait_for_channel_received(collector, channel_id, timeout_s, poll_interval_s=
time.sleep(poll_interval_s)
def wait_for_channel_received_count(collector, channel_id, count, timeout_s, poll_interval_s=0.5):
"""Channel events for `channel_id`, oldest first; waits for `count` of them.
Returns whatever has arrived when the timeout expires, so a caller can
assert on both the order and the shortfall.
"""
deadline = time.monotonic() + timeout_s
while True:
received = [e for e in collector.snapshot() if e.get("eventType") == CHANNEL_RECEIVED_EVENT and e.get("channelId") == channel_id]
if len(received) >= count or time.monotonic() >= deadline:
return received
time.sleep(poll_interval_s)
def channel_payloads(events):
"""Decoded payloads of channel events, oldest first."""
return [base64.b64decode(event["payload"]) for event in events]
def wait_for_message_received(collector, content_topic, timeout_s, poll_interval_s=0.5):
"""Wait for a messaging-layer message_received event on `content_topic`.
@@ -186,6 +229,324 @@ class TestChannelDelivery:
leaked = wait_for_channel_received(receiver_collector, channel_id, NO_CHANNEL_DELIVERY_WINDOW_S)
assert leaked is None, f"an unmarked message must not surface as a {CHANNEL_RECEIVED_EVENT}; got: {leaked!r}"
def test_rc07_wrong_content_topic_dropped(self, node_config):
"""RC07: a correctly-marked channel message published on a different
content topic is dropped at channel ingress.
A runs the same channelId as B but on RC07_OTHER_CONTENT_TOPIC. B is
subscribed to that topic too, so B's messaging layer still surfaces the
message (message_received), proving it arrived, but B's channel is bound
to RC07_CONTENT_TOPIC and must not fire channel_message_received.
"""
channel_id = unique_channel_id(RC07_CHANNEL_PREFIX)
payload_b64 = to_base64("rc07 message on the wrong content topic")
node_config.update(
{
"relay": True,
"store": False,
"reliabilityEnabled": False,
"numShardsInNetwork": 1,
}
)
receiver_collector = EventCollector()
receiver_result = WrapperManager.create_and_start(config=node_config, event_cb=receiver_collector.event_callback)
assert receiver_result.is_ok(), f"Failed to start receiver: {receiver_result.err()}"
with receiver_result.ok_value as receiver:
sender_config = {
**node_config,
"staticnodes": [get_node_multiaddr(receiver)],
"portsShift": 1,
}
for content_topic in (RC07_CONTENT_TOPIC, RC07_OTHER_CONTENT_TOPIC):
subscribe_result = receiver.subscribe_content_topic(content_topic)
assert subscribe_result.is_ok(), f"receiver subscribe_content_topic {content_topic} failed: {subscribe_result.err()}"
receiver_create = receiver.channel_create(channel_id, RC07_CONTENT_TOPIC, SENDER_B)
assert receiver_create.is_ok(), f"receiver channel_create failed: {receiver_create.err()}"
with ChannelSenderProcess(
sender_config,
content_topic=RC07_OTHER_CONTENT_TOPIC,
channel_id=channel_id,
sender_id=SENDER_A,
payload_b64=payload_b64,
settle_s=MESH_SETTLE_S,
):
arrived = wait_for_message_received(receiver_collector, RC07_OTHER_CONTENT_TOPIC, DELIVERY_TIMEOUT_S)
assert arrived is not None, (
f"receiver never saw a {MESSAGE_RECEIVED_EVENT} for {RC07_OTHER_CONTENT_TOPIC} within {DELIVERY_TIMEOUT_S}s; "
f"cannot conclude the channel dropped it. Collected events: {receiver_collector.snapshot()}"
)
leaked = wait_for_channel_received(receiver_collector, channel_id, NO_CHANNEL_DELIVERY_WINDOW_S)
assert leaked is None, f"a message on a foreign content topic must not surface as a {CHANNEL_RECEIVED_EVENT}; got: {leaked!r}"
def test_rc08_bidirectional_exchange_preserves_causal_order(self, node_config):
"""RC08: A and B exchange interleaved messages on one channel, and each
side delivers the other's in causal order.
Each reply is sent only after the message it answers has landed, so m2
references m1 and m3 references m2. B must deliver m1 then m3, A must
deliver m2.
"""
channel_id = unique_channel_id(RC08_CHANNEL_PREFIX)
m1, m2, m3 = "rc08 from A", "rc08 reply from B", "rc08 follow-up from A"
node_config.update(
{
"relay": True,
"store": False,
"reliabilityEnabled": False,
"numShardsInNetwork": 1,
}
)
receiver_collector = EventCollector()
receiver_result = WrapperManager.create_and_start(config=node_config, event_cb=receiver_collector.event_callback)
assert receiver_result.is_ok(), f"Failed to start receiver: {receiver_result.err()}"
with receiver_result.ok_value as receiver:
sender_config = {
**node_config,
"staticnodes": [get_node_multiaddr(receiver)],
"portsShift": 1,
}
subscribe_result = receiver.subscribe_content_topic(RC08_CONTENT_TOPIC)
assert subscribe_result.is_ok(), f"receiver subscribe_content_topic failed: {subscribe_result.err()}"
receiver_create = receiver.channel_create(channel_id, RC08_CONTENT_TOPIC, SENDER_B)
assert receiver_create.is_ok(), f"receiver channel_create failed: {receiver_create.err()}"
with ChannelSenderProcess(
sender_config,
content_topic=RC08_CONTENT_TOPIC,
channel_id=channel_id,
sender_id=SENDER_A,
payload_b64=to_base64(m1),
settle_s=MESH_SETTLE_S,
) as sender:
first = wait_for_channel_received(receiver_collector, channel_id, DELIVERY_TIMEOUT_S)
assert first is not None, (
f"B never saw A's first message within {DELIVERY_TIMEOUT_S}s; "
f"nothing to reply to. Collected events: {receiver_collector.snapshot()}"
)
# B replies only now, so m2 carries m1 in its causal history.
delay(RC08_CAUSAL_SETTLE_S)
reply_result = receiver.channel_send(channel_id, create_message_bindings(payload=to_base64(m2)))
assert reply_result.is_ok(), f"receiver channel_send failed: {reply_result.err()}"
delivered_to_a = sender.wait_for_received(1, DELIVERY_TIMEOUT_S)
assert channel_payloads(delivered_to_a) == [m2.encode()], f"A must deliver B's reply; got {channel_payloads(delivered_to_a)!r}"
# Same race on A's side.
delay(RC08_CAUSAL_SETTLE_S)
sender.send(to_base64(m3))
delivered_to_b = wait_for_channel_received_count(receiver_collector, channel_id, 2, DELIVERY_TIMEOUT_S)
assert channel_payloads(delivered_to_b) == [m1.encode(), m3.encode()], (
f"B must deliver A's messages in causal order; got {channel_payloads(delivered_to_b)!r}. "
f"Collected events: {receiver_collector.snapshot()}"
)
def test_rc09_late_joining_receiver_still_receives(self, node_config):
"""RC09: A sends m1 while B is not up yet; B joins and still ends up with
both m1 and m2.
Recovery here is the send service's relay retry, not SDS: A republishes
m1 while it is unacknowledged (within MaxTimeInCache) and it lands as
soon as B joins the mesh. SDS-R is not exercised — a repair request only
rides out after T_req, a hash in [repairTMin, repairTMax) = [30s, 300s)
that logos-delivery does not expose, so it cannot be driven from an E2E
test. That path is RC10 — see test_rc10_missing_dependency_is_parked
below and tests/wrappers_tests/test_channel_repair.py.
"""
channel_id = unique_channel_id(RC09_CHANNEL_PREFIX)
m1, m2 = "rc09 sent while B is away", "rc09 sent after B joins"
node_config.update(
{
"relay": True,
"store": False,
"reliabilityEnabled": False,
"numShardsInNetwork": 1,
}
)
# A comes up first with nobody to dial; B joins later and dials A.
sender_config = {**node_config, "portsShift": 1}
with ChannelSenderProcess(
sender_config,
content_topic=RC09_CONTENT_TOPIC,
channel_id=channel_id,
sender_id=SENDER_A,
payload_b64=to_base64(m1),
settle_s=RC09_SENDER_SETTLE_S,
) as sender:
receiver_collector = EventCollector()
receiver_config = {**node_config, "staticnodes": [sender.multiaddr]}
receiver_result = WrapperManager.create_and_start(config=receiver_config, event_cb=receiver_collector.event_callback)
assert receiver_result.is_ok(), f"Failed to start receiver: {receiver_result.err()}"
with receiver_result.ok_value as receiver:
assert wait_for_connected(receiver_collector) is not None, "Receiver did not reach Connected/PartiallyConnected state"
subscribe_result = receiver.subscribe_content_topic(RC09_CONTENT_TOPIC)
assert subscribe_result.is_ok(), f"receiver subscribe_content_topic failed: {subscribe_result.err()}"
receiver_create = receiver.channel_create(channel_id, RC09_CONTENT_TOPIC, SENDER_B)
assert receiver_create.is_ok(), f"receiver channel_create failed: {receiver_create.err()}"
delay(MESH_SETTLE_S)
sender.send(to_base64(m2))
recovered = wait_for_channel_received_count(receiver_collector, channel_id, 2, RC09_RECOVERY_TIMEOUT_S)
assert channel_payloads(recovered) == [m1.encode(), m2.encode()], (
f"B must end up with the message sent before it joined, then the one after; "
f"got {channel_payloads(recovered)!r}. Collected events: {receiver_collector.snapshot()}"
)
def test_rc12_three_participants_all_receive(self, node_config):
"""RC12: A, B and C share one channel; A's send reaches both B and C.
Each delivers m1 exactly once, with A's senderId and the exact payload.
B is the node under test; A and C run as storage-isolated peers dialing it.
"""
channel_id = unique_channel_id(RC12_CHANNEL_PREFIX)
m1 = "rc12 from A"
node_config.update(
{
"relay": True,
"store": False,
"reliabilityEnabled": False,
"numShardsInNetwork": 1,
}
)
receiver_collector = EventCollector()
receiver_result = WrapperManager.create_and_start(config=node_config, event_cb=receiver_collector.event_callback)
assert receiver_result.is_ok(), f"Failed to start receiver: {receiver_result.err()}"
with receiver_result.ok_value as receiver:
peer_config = {**node_config, "staticnodes": [get_node_multiaddr(receiver)]}
subscribe_result = receiver.subscribe_content_topic(RC12_CONTENT_TOPIC)
assert subscribe_result.is_ok(), f"receiver subscribe_content_topic failed: {subscribe_result.err()}"
receiver_create = receiver.channel_create(channel_id, RC12_CONTENT_TOPIC, SENDER_B)
assert receiver_create.is_ok(), f"receiver channel_create failed: {receiver_create.err()}"
# C joins first so it is meshed by the time A sends m1.
with ChannelSenderProcess(
{**peer_config, "portsShift": 2},
content_topic=RC12_CONTENT_TOPIC,
channel_id=channel_id,
sender_id=SENDER_C,
settle_s=MESH_SETTLE_S,
) as third_party:
with ChannelSenderProcess(
{**peer_config, "portsShift": 1},
content_topic=RC12_CONTENT_TOPIC,
channel_id=channel_id,
sender_id=SENDER_A,
payload_b64=to_base64(m1),
settle_s=MESH_SETTLE_S,
):
assert (
wait_for_channel_received(receiver_collector, channel_id, DELIVERY_TIMEOUT_S) is not None
), f"B never saw A's message within {DELIVERY_TIMEOUT_S}s. Collected events: {receiver_collector.snapshot()}"
assert third_party.wait_for_received(1, DELIVERY_TIMEOUT_S), f"C never saw A's message within {DELIVERY_TIMEOUT_S}s"
# B and C must have exactly one delivery; any second event is a duplicate.
on_b = wait_for_channel_received_count(receiver_collector, channel_id, 2, NO_CHANNEL_DELIVERY_WINDOW_S)
on_c = third_party.wait_for_received(2, NO_CHANNEL_DELIVERY_WINDOW_S)
assert channel_payloads(on_b) == [m1.encode()], f"B must deliver A's message exactly once; got {channel_payloads(on_b)!r}"
assert channel_payloads(on_c) == [m1.encode()], f"C must deliver A's message exactly once; got {channel_payloads(on_c)!r}"
assert [e["senderId"] for e in on_b + on_c] == [
SENDER_A,
SENDER_A,
], f"both events must carry {SENDER_A!r}, got {[e.get('senderId') for e in on_b + on_c]!r}"
def test_rc13_close_recreate_then_send_delivers_new_message(self, node_config):
"""RC13: A closes and re-creates its channel, then sends again.
A sends m1, cycles the channel under the same id, then sends m2. B must
deliver m2 exactly once and ordered after m1, and must not replay m1 —
the re-created channel picks up the restored SDS history rather than
starting a fresh one.
Distinct from the nim in-process test, which cycles the *receiver's*
channel and asserts a replayed m1 is suppressed on ingress; here the
cycle is on the send path.
"""
channel_id = unique_channel_id(RC13_CHANNEL_PREFIX)
m1, m2 = "rc13 before close", "rc13 after re-create"
node_config.update(
{
"relay": True,
"store": False,
"reliabilityEnabled": False,
"numShardsInNetwork": 1,
}
)
receiver_collector = EventCollector()
receiver_result = WrapperManager.create_and_start(config=node_config, event_cb=receiver_collector.event_callback)
assert receiver_result.is_ok(), f"Failed to start receiver: {receiver_result.err()}"
with receiver_result.ok_value as receiver:
sender_config = {
**node_config,
"staticnodes": [get_node_multiaddr(receiver)],
"portsShift": 1,
}
subscribe_result = receiver.subscribe_content_topic(RC13_CONTENT_TOPIC)
assert subscribe_result.is_ok(), f"receiver subscribe_content_topic failed: {subscribe_result.err()}"
receiver_create = receiver.channel_create(channel_id, RC13_CONTENT_TOPIC, SENDER_B)
assert receiver_create.is_ok(), f"receiver channel_create failed: {receiver_create.err()}"
with ChannelSenderProcess(
sender_config,
content_topic=RC13_CONTENT_TOPIC,
channel_id=channel_id,
sender_id=SENDER_A,
payload_b64=to_base64(m1),
settle_s=MESH_SETTLE_S,
) as sender:
first = wait_for_channel_received(receiver_collector, channel_id, DELIVERY_TIMEOUT_S)
assert first is not None, (
f"No {CHANNEL_RECEIVED_EVENT} for m1 on {channel_id} within {DELIVERY_TIMEOUT_S}s; "
f"the close/re-create is only meaningful once m1 landed. Collected events: {receiver_collector.snapshot()}"
)
sender.close_and_recreate()
sender.send(to_base64(m2))
received = wait_for_channel_received_count(receiver_collector, channel_id, 2, DELIVERY_TIMEOUT_S)
assert len(received) == 2, (
f"expected m1 then m2 on {channel_id}, got {len(received)}: {channel_payloads(received)!r}. "
f"Collected events: {receiver_collector.snapshot()}"
)
assert channel_payloads(received) == [
m1.encode(),
m2.encode(),
], f"expected [m1, m2] in causal order, got: {channel_payloads(received)!r}"
# A third event could only be a replayed m1 from the re-created channel.
settled = wait_for_channel_received_count(receiver_collector, channel_id, 3, NO_CHANNEL_DELIVERY_WINDOW_S)
assert len(settled) == 2, f"re-create must not re-deliver m1; got: {channel_payloads(settled)!r}"
def test_receive_after_close_emits_no_channel_event(self, node_config):
"""A closed channel must not deliver: B closes its channel, then A sends a
valid channel message on the same channel + content topic.
@@ -14,6 +14,7 @@ RC03_OPEN_CHANNEL_PREFIX = "rc03-open"
RC03_RANDOM_CHANNEL_PREFIX = "rc03-random"
RC04_CHANNEL_PREFIX = "rc04-channel"
RC04_CLOSED_CHANNEL_PREFIX = "rc04-closed-channel"
RC04_EPHEMERAL_CHANNEL_PREFIX = "rc04-ephemeral-channel"
RC04_DISTINCT_CHANNEL_PREFIX = "rc04-distinct-channel"
@@ -176,6 +177,38 @@ class TestChannelLifecycle:
finally:
node.stop_and_destroy()
def test_rc04_send_after_close_rejected(self, node_config):
"""RC04 variant: sending after the channel is closed is rejected.
Exercises the teardown -> send ordering (distinct from RC02's
never-created id): once channel_close removes the id, channel_send Errs
with "unknown channel: <id>". No events are expected.
"""
node_config.update({"mode": "Core", "numShardsInNetwork": 1})
channel_id = unique_channel_id(RC04_CLOSED_CHANNEL_PREFIX)
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)
close_result = node.channel_close(channel_id)
assert close_result.is_ok(), f"channel_close failed: {close_result.err()}"
send_result = node.channel_send(channel_id, create_message_bindings())
assert send_result.is_err(), f"channel_send after close must fail, got Ok({send_result.ok_value!r})"
assert f"unknown channel: {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_rc04_ephemeral_send_returns_handle(self, node_config):
"""RC04 variant: an ephemeral channel_send still returns a handle synchronously.