mirror of
https://github.com/varun-r-mallya/py-libp2p.git
synced 2025-12-31 20:36:24 +00:00
feat: add support for sparse connect (#680)
* init * add newsfragment * fix
This commit is contained in:
@ -17,6 +17,7 @@ from tests.utils.factories import (
|
||||
from tests.utils.pubsub.utils import (
|
||||
dense_connect,
|
||||
one_to_all_connect,
|
||||
sparse_connect,
|
||||
)
|
||||
|
||||
|
||||
@ -506,3 +507,84 @@ async def test_gossip_heartbeat(initial_peer_count, monkeypatch):
|
||||
# Check that the peer to gossip to is not in our fanout peers
|
||||
assert peer not in fanout_peers
|
||||
assert topic_fanout in peers_to_gossip[peer]
|
||||
|
||||
|
||||
@pytest.mark.trio
|
||||
async def test_dense_connect_fallback():
|
||||
"""Test that sparse_connect falls back to dense connect for small networks."""
|
||||
async with PubsubFactory.create_batch_with_gossipsub(3) as pubsubs_gsub:
|
||||
hosts = [pubsub.host for pubsub in pubsubs_gsub]
|
||||
degree = 2
|
||||
|
||||
# Create network (should use dense connect)
|
||||
await sparse_connect(hosts, degree)
|
||||
|
||||
# Wait for connections to be established
|
||||
await trio.sleep(2)
|
||||
|
||||
# Verify dense topology (all nodes connected to each other)
|
||||
for i, pubsub in enumerate(pubsubs_gsub):
|
||||
connected_peers = len(pubsub.peers)
|
||||
expected_connections = len(hosts) - 1
|
||||
assert connected_peers == expected_connections, (
|
||||
f"Host {i} has {connected_peers} connections, "
|
||||
f"expected {expected_connections} in dense mode"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.trio
|
||||
async def test_sparse_connect():
|
||||
"""Test sparse connect functionality and message propagation."""
|
||||
async with PubsubFactory.create_batch_with_gossipsub(10) as pubsubs_gsub:
|
||||
hosts = [pubsub.host for pubsub in pubsubs_gsub]
|
||||
degree = 2
|
||||
topic = "test_topic"
|
||||
|
||||
# Create network (should use sparse connect)
|
||||
await sparse_connect(hosts, degree)
|
||||
|
||||
# Wait for connections to be established
|
||||
await trio.sleep(2)
|
||||
|
||||
# Verify sparse topology
|
||||
for i, pubsub in enumerate(pubsubs_gsub):
|
||||
connected_peers = len(pubsub.peers)
|
||||
assert degree <= connected_peers < len(hosts) - 1, (
|
||||
f"Host {i} has {connected_peers} connections, "
|
||||
f"expected between {degree} and {len(hosts) - 1} in sparse mode"
|
||||
)
|
||||
|
||||
# Test message propagation
|
||||
queues = [await pubsub.subscribe(topic) for pubsub in pubsubs_gsub]
|
||||
await trio.sleep(2)
|
||||
|
||||
# Publish and verify message propagation
|
||||
msg_content = b"test_msg"
|
||||
await pubsubs_gsub[0].publish(topic, msg_content)
|
||||
await trio.sleep(2)
|
||||
|
||||
# Verify message propagation - ideally all nodes should receive it
|
||||
received_count = 0
|
||||
for queue in queues:
|
||||
try:
|
||||
msg = await queue.get()
|
||||
if msg.data == msg_content:
|
||||
received_count += 1
|
||||
except Exception:
|
||||
continue
|
||||
|
||||
total_nodes = len(pubsubs_gsub)
|
||||
|
||||
# Ideally all nodes should receive the message for optimal scalability
|
||||
if received_count == total_nodes:
|
||||
# Perfect propagation achieved
|
||||
pass
|
||||
else:
|
||||
# require more than half for acceptable scalability
|
||||
min_required = (total_nodes + 1) // 2
|
||||
assert received_count >= min_required, (
|
||||
f"Message propagation insufficient: "
|
||||
f"{received_count}/{total_nodes} nodes "
|
||||
f"received the message. Ideally all nodes should receive it, but at "
|
||||
f"minimum {min_required} required for sparse network scalability."
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user