Files
py-libp2p/libp2p/transport/tcp/tcp.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

193 lines
6.1 KiB
Python

from collections.abc import (
Awaitable,
Callable,
Sequence,
)
import logging
from multiaddr import (
Multiaddr,
)
import trio
from trio_typing import (
TaskStatus,
)
from libp2p.abc import (
IListener,
IRawConnection,
ITransport,
)
from libp2p.custom_types import (
THandler,
)
from libp2p.io.trio import (
TrioTCPStream,
)
from libp2p.network.connection.raw_connection import (
RawConnection,
)
from libp2p.transport.exceptions import (
OpenConnectionError,
)
logger = logging.getLogger("libp2p.transport.tcp")
class TCPListener(IListener):
listeners: list[trio.SocketListener]
def __init__(self, handler_function: THandler) -> None:
self.listeners = []
self.handler = handler_function
# TODO: Get rid of `nursery`?
async def listen(self, maddr: Multiaddr, nursery: trio.Nursery) -> bool:
"""
Put listener in listening mode and wait for incoming connections.
:param maddr: maddr of peer
:return: return True if successful
"""
async def serve_tcp(
handler: Callable[[trio.SocketStream], Awaitable[None]],
port: int,
host: str,
task_status: TaskStatus[Sequence[trio.SocketListener]],
) -> None:
"""Just a proxy function to add logging here."""
logger.debug("serve_tcp %s %s", host, port)
await trio.serve_tcp(handler, port, host=host, task_status=task_status)
async def handler(stream: trio.SocketStream) -> None:
remote_host: str = ""
remote_port: int = 0
try:
tcp_stream = TrioTCPStream(stream)
remote_tuple = tcp_stream.get_remote_address()
if remote_tuple is not None:
remote_host, remote_port = remote_tuple
await self.handler(tcp_stream)
except Exception:
logger.debug(f"Connection from {remote_host}:{remote_port} failed.")
tcp_port_str = maddr.value_for_protocol("tcp")
if tcp_port_str is None:
logger.error(f"Cannot listen: TCP port is missing in multiaddress {maddr}")
return False
try:
tcp_port = int(tcp_port_str)
except ValueError:
logger.error(
f"Cannot listen: Invalid TCP port '{tcp_port_str}' "
f"in multiaddress {maddr}"
)
return False
ip4_host_str = maddr.value_for_protocol("ip4")
# For trio.serve_tcp, ip4_host_str (as host argument) can be None,
# which typically means listen on all available interfaces.
started_listeners = await nursery.start(
serve_tcp,
handler,
tcp_port,
ip4_host_str,
)
if started_listeners is None:
# This implies that task_status.started() was not called within serve_tcp,
# likely because trio.serve_tcp itself failed to start (e.g., port in use).
logger.error(
f"Failed to start TCP listener for {maddr}: "
f"`nursery.start` returned None. "
"This might be due to issues like the port already "
"being in use or invalid host."
)
return False
self.listeners.extend(started_listeners)
return True
def get_addrs(self) -> tuple[Multiaddr, ...]:
"""
Retrieve list of addresses the listener is listening on.
:return: return list of addrs
"""
return tuple(
_multiaddr_from_socket(listener.socket) for listener in self.listeners
)
async def close(self) -> None:
async with trio.open_nursery() as nursery:
for listener in self.listeners:
nursery.start_soon(listener.aclose)
class TCP(ITransport):
async def dial(self, maddr: Multiaddr) -> IRawConnection:
"""
Dial a transport to peer listening on multiaddr.
:param maddr: multiaddr of peer
:return: `RawConnection` if successful
:raise OpenConnectionError: raised when failed to open connection
"""
host_str = maddr.value_for_protocol("ip4")
port_str = maddr.value_for_protocol("tcp")
if host_str is None:
raise OpenConnectionError(
f"Failed to dial {maddr}: IP address not found in multiaddr."
)
if port_str is None:
raise OpenConnectionError(
f"Failed to dial {maddr}: TCP port not found in multiaddr."
)
try:
port_int = int(port_str)
except ValueError:
raise OpenConnectionError(
f"Failed to dial {maddr}: Invalid TCP port '{port_str}'."
)
try:
# trio.open_tcp_stream requires host to be str or bytes, not None.
stream = await trio.open_tcp_stream(host_str, port_int)
except OSError as error:
# OSError is common for network issues like "Connection refused"
# or "Host unreachable".
raise OpenConnectionError(
f"Failed to open TCP stream to {maddr}: {error}"
) from error
except Exception as error:
# Catch other potential errors from trio.open_tcp_stream and wrap them.
raise OpenConnectionError(
f"An unexpected error occurred when dialing {maddr}: {error}"
) from error
read_write_closer = TrioTCPStream(stream)
return RawConnection(read_write_closer, True)
def create_listener(self, handler_function: THandler) -> TCPListener:
"""
Create listener on transport.
:param handler_function: a function called when a new connection is received
that takes a connection as argument which implements interface-connection
:return: a listener object that implements listener_interface.py
"""
return TCPListener(handler_function)
def _multiaddr_from_socket(socket: trio.socket.SocketType) -> Multiaddr:
ip, port = socket.getsockname()
return Multiaddr(f"/ip4/{ip}/tcp/{port}")