Files
py-libp2p/libp2p/tools/utils.py
Arush Kurundodi bdadec7519 ft. modernise py-libp2p (#618)
* fix pyproject.toml , add ruff

* rm lock

* make progress

* add poetry lock ignore

* fix type issues

* fix tcp type errors

* fix text example - type error - wrong args

* add setuptools to dev

* test ci

* fix docs build

* fix type issues for new_swarm & new_host

* fix types in gossipsub

* fix type issues in noise

* wip: factories

* revert factories

* fix more type issues

* more type fixes

* fix: add null checks for noise protocol initialization and key handling

* corrected argument-errors in peerId and Multiaddr in peer tests

* fix: Noice - remove redundant type casts in BaseNoiseMsgReadWriter

* fix: update test_notify.py to use SwarmFactory.create_batch_and_listen, fix type hints, and comment out ClosedStream assertions

* Fix type checks for pubsub module

Signed-off-by: sukhman <sukhmansinghsaluja@gmail.com>

* Fix type checks for pubsub module-tests

Signed-off-by: sukhman <sukhmansinghsaluja@gmail.com>

* noise: add checks for uninitialized protocol and key states in PatternXX

Signed-off-by: varun-r-mallya <varunrmallya@gmail.com>

* pubsub: add None checks for optional fields in FloodSub and Pubsub

Signed-off-by: varun-r-mallya <varunrmallya@gmail.com>

* Fix type hints and improve testing

Signed-off-by: varun-r-mallya <varunrmallya@gmail.com>

* remove redundant checks

Signed-off-by: varun-r-mallya <varunrmallya@gmail.com>

* fix build issues

* add optional to trio service

* fix types

* fix type errors

* Fix type errors

Signed-off-by: varun-r-mallya <varunrmallya@gmail.com>

* fixed more-type checks in crypto and peer_data files

* wip: factories

* replaced union with optional

* fix: type-error in interp-utils and peerinfo

* replace pyright with pyrefly

* add pyrefly.toml

* wip: fix multiselect issues

* try typecheck

* base check

* mcache test fixes , typecheck ci update

* fix ci

* will this work

* minor fix

* use poetry

* fix wokflow

* use cache,fix err

* fix pyrefly.toml

* fix pyrefly.toml

* fix cache in ci

* deploy commit

* add main baseline

* update to v5

* improve typecheck ci (#14)

* fix typo

* remove holepunching code (#16)

* fix gossipsub typeerrors (#17)

* fix: ensure initiator user includes remote peer id in handshake (#15)

* fix ci (#19)

* typefix: custom_types | core/peerinfo/test_peer_info | io/abc | pubsub/floodsub | protocol_muxer/multiselect (#18)

* fix: Typefixes in PeerInfo  (#21)

* fix minor type issue (#22)

* fix type errors in pubsub (#24)

* fix: Minor typefixes in tests (#23)

* Fix failing tests for type-fixed test/pubsub (#8)

* move pyrefly & ruff to pyproject.toml & rm .project-template (#28)

* move the async_context file to tests/core

* move crypto test to crypto folder

* fix: some typefixes (#25)

* fix type errors

* fix type issues

* fix: update gRPC API usage in autonat_pb2_grpc.py (#31)

* md: typecheck ci

* rm comments

* clean up : from review suggestions

* use | None over Optional as per new python standards

* drop supporto for py3.9

* newsfragments

---------

Signed-off-by: sukhman <sukhmansinghsaluja@gmail.com>
Signed-off-by: varun-r-mallya <varunrmallya@gmail.com>
Co-authored-by: acul71 <luca.pisani@birdo.net>
Co-authored-by: kaneki003 <sakshamchauhan707@gmail.com>
Co-authored-by: sukhman <sukhmansinghsaluja@gmail.com>
Co-authored-by: varun-r-mallya <varunrmallya@gmail.com>
Co-authored-by: varunrmallya <100590632+varun-r-mallya@users.noreply.github.com>
Co-authored-by: lla-dane <abhinavagarwalla6@gmail.com>
Co-authored-by: Collins <ArtemisfowlX@protonmail.com>
Co-authored-by: Abhinav Agarwalla <120122716+lla-dane@users.noreply.github.com>
Co-authored-by: guha-rahul <52607971+guha-rahul@users.noreply.github.com>
Co-authored-by: Sukhman Singh <63765293+sukhman-sukh@users.noreply.github.com>
Co-authored-by: acul71 <34693171+acul71@users.noreply.github.com>
Co-authored-by: pacrob <5199899+pacrob@users.noreply.github.com>
2025-06-09 11:39:59 -06:00

153 lines
4.1 KiB
Python

from collections.abc import (
Awaitable,
Callable,
)
import logging
import trio
from libp2p.abc import (
IHost,
INetStream,
)
from libp2p.network.stream.exceptions import (
StreamError,
)
from libp2p.network.swarm import (
Swarm,
)
from libp2p.peer.peerinfo import (
info_from_p2p_addr,
)
from .constants import (
MAX_READ_LEN,
)
async def connect_swarm(swarm_0: Swarm, swarm_1: Swarm) -> None:
peer_id = swarm_1.get_peer_id()
addrs = tuple(
addr
for transport in swarm_1.listeners.values()
for addr in transport.get_addrs()
)
swarm_0.peerstore.add_addrs(peer_id, addrs, 10000)
# Add retry logic for more robust connection
max_retries = 3
retry_delay = 0.2
last_error = None
for attempt in range(max_retries):
try:
await swarm_0.dial_peer(peer_id)
# Verify connection is established in both directions
if (
swarm_0.get_peer_id() in swarm_1.connections
and swarm_1.get_peer_id() in swarm_0.connections
):
return
# Connection partially established, wait a bit for it to complete
await trio.sleep(0.1)
if (
swarm_0.get_peer_id() in swarm_1.connections
and swarm_1.get_peer_id() in swarm_0.connections
):
return
logging.debug(
"Swarm connection verification failed on attempt"
+ f" {attempt + 1}, retrying..."
)
except Exception as e:
last_error = e
logging.debug(f"Swarm connection attempt {attempt + 1} failed: {e}")
await trio.sleep(retry_delay)
# If we got here, all retries failed
if last_error:
raise RuntimeError(
f"Failed to connect swarms after {max_retries} attempts"
) from last_error
else:
err_msg = (
"Failed to establish bidirectional swarm connection"
+ f" after {max_retries} attempts"
)
raise RuntimeError(err_msg)
async def connect(node1: IHost, node2: IHost) -> None:
"""Connect node1 to node2."""
addr = node2.get_addrs()[0]
info = info_from_p2p_addr(addr)
# Add retry logic for more robust connection
max_retries = 3
retry_delay = 0.2
last_error = None
for attempt in range(max_retries):
try:
await node1.connect(info)
# Verify connection is established in both directions
if (
node2.get_id() in node1.get_network().connections
and node1.get_id() in node2.get_network().connections
):
return
# Connection partially established, wait a bit for it to complete
await trio.sleep(0.1)
if (
node2.get_id() in node1.get_network().connections
and node1.get_id() in node2.get_network().connections
):
return
logging.debug(
f"Connection verification failed on attempt {attempt + 1}, retrying..."
)
except Exception as e:
last_error = e
logging.debug(f"Connection attempt {attempt + 1} failed: {e}")
await trio.sleep(retry_delay)
# If we got here, all retries failed
if last_error:
raise RuntimeError(
f"Failed to connect after {max_retries} attempts"
) from last_error
else:
err_msg = (
f"Failed to establish bidirectional connection after {max_retries} attempts"
)
raise RuntimeError(err_msg)
def create_echo_stream_handler(
ack_prefix: str,
) -> Callable[[INetStream], Awaitable[None]]:
async def echo_stream_handler(stream: INetStream) -> None:
while True:
try:
read_string = (await stream.read(MAX_READ_LEN)).decode()
except StreamError:
break
resp = ack_prefix + read_string
try:
await stream.write(resp.encode())
except StreamError:
break
return echo_stream_handler