Skip to content

proxystore.endpoint.p2p.manager

Manager of peer connections to other endpoints.

RequestHandler module-attribute

RequestHandler = Callable[
    [EndpointId, Message], Awaitable[Message]
]

Handler of requests from peers.

The handler is called with the ID of the peer and the request and returns the response.

PathInfo dataclass

PathInfo(relayed: bool, remote_addr: str, rtt_ms: int)

Network path used by a connection to a peer.

Attributes:

  • relayed (bool) –

    If traffic is relayed rather than sent directly to the peer.

  • remote_addr (str) –

    Address of the peer (or relay) on this path.

  • rtt_ms (int) –

    Round-trip time in milliseconds estimated by QUIC.

from_connection classmethod

from_connection(connection: Connection) -> Self | None

Get the path selected for sending on a connection.

Returns:

  • Self | None –

    The selected path or None if the connection has no path.

Source code in proxystore/endpoint/p2p/manager.py
@classmethod
def from_connection(cls, connection: iroh.Connection) -> Self | None:
    """Get the path selected for sending on a connection.

    Returns:
        The selected path or `None` if the connection has no path.
    """
    for path in connection.paths():
        if path.is_selected:
            return cls(
                relayed=path.is_relay,
                remote_addr=path.remote_addr,
                rtt_ms=path.rtt_ms,
            )
    return None

describe

describe() -> str

Describe the path for logging.

Source code in proxystore/endpoint/p2p/manager.py
def describe(self) -> str:
    """Describe the path for logging."""
    kind = 'relayed via' if self.relayed else 'direct to'
    return f'{kind} {self.remote_addr} (rtt {self.rtt_ms} ms)'

PeerConnection dataclass

PeerConnection(
    peer_id: EndpointId,
    connection: Connection,
    dialed: bool,
)

Connection to a peer.

Attributes:

  • peer_id (EndpointId) –

    ID of the peer.

  • connection (Connection) –

    Underlying iroh connection.

  • dialed (bool) –

    If this endpoint opened the connection (rather than the peer).

closed property

closed: bool

The connection is closed.

version property

version: int

Protocol version negotiated on the connection.

path

path() -> PathInfo | None

Get the path selected for sending on the connection.

Source code in proxystore/endpoint/p2p/manager.py
def path(self) -> PathInfo | None:
    """Get the path selected for sending on the connection."""
    return PathInfo.from_connection(self.connection)

report_path

report_path(peer_name: str) -> None

Log the path of the connection if it changed since last reported.

Changes in only the RTT of the path are not reported.

Parameters:

  • peer_name (str) –

    Name of the peer for the log.

Source code in proxystore/endpoint/p2p/manager.py
def report_path(self, peer_name: str) -> None:
    """Log the path of the connection if it changed since last reported.

    Changes in only the RTT of the path are not reported.

    Args:
        peer_name: Name of the peer for the log.
    """
    path = self.path()
    key = None if path is None else (path.relayed, path.remote_addr)
    if self._reported and key == self._reported_path:
        return
    self._reported = True
    self._reported_path = key
    direction = 'to' if self.dialed else 'from'
    if path is None:
        logger.info(
            'Connection %s peer %s has no path',
            direction,
            peer_name,
        )
    else:
        logger.info(
            'Connection %s peer %s is %s',
            direction,
            peer_name,
            path.describe(),
        )

watch_path async

watch_path(peer_name: str) -> None

Report changes to the path of the connection after connecting.

A new connection often starts relayed and switches to a direct path once hole-punching succeeds, so the path is checked periodically for a short time to report the change even if the connection is idle.

Parameters:

  • peer_name (str) –

    Name of the peer for the log.

Source code in proxystore/endpoint/p2p/manager.py
async def watch_path(self, peer_name: str) -> None:
    """Report changes to the path of the connection after connecting.

    A new connection often starts relayed and switches to a direct path
    once hole-punching succeeds, so the path is checked periodically for
    a short time to report the change even if the connection is idle.

    Args:
        peer_name: Name of the peer for the log.
    """
    loop = asyncio.get_running_loop()
    end = loop.time() + _PATH_WATCH_DURATION
    while loop.time() < end:
        await asyncio.sleep(_PATH_WATCH_INTERVAL)
        if self.closed:
            return
        self.report_path(peer_name)

PeerPolicy

Bases: Protocol

Policy of which peer endpoints an endpoint communicates with.

The PeerManager checks the policy on each connection and request, and closes the connections to peers which are no longer allowed. The Allowlist is the policy of endpoints started from an endpoint directory.

allowed

allowed(peer_id: EndpointId) -> bool

Check if the endpoint is allowed to communicate with this one.

Source code in proxystore/endpoint/p2p/manager.py
def allowed(self, peer_id: EndpointId) -> bool:
    """Check if the endpoint is allowed to communicate with this one."""
    ...

name_of

name_of(peer_id: EndpointId) -> str | None

Get the name of the peer used in logs or None if unknown.

Source code in proxystore/endpoint/p2p/manager.py
def name_of(self, peer_id: EndpointId) -> str | None:
    """Get the name of the peer used in logs or `None` if unknown."""
    ...

PeerOptions dataclass

PeerOptions(
    preset: Preset = preset_n0(),
    relay_mode: RelayMode | None = None,
    bind_addr: str | None = None,
    connect_timeout: float = 30,
    online_timeout: float | None = 10,
)

Options of the connections of a peer manager to peers.

Attributes:

  • preset (Preset) –

    iroh preset used to configure discovery and relays. Defaults to iroh.preset_n0() which uses n0's public relays and DNS discovery.

  • relay_mode (RelayMode | None) –

    Relay mode which overrides the relays of the preset or None to use the relays of the preset.

  • bind_addr (str | None) –

    Address to bind to (e.g., "127.0.0.1:0") or None to bind to all interfaces on a random port.

  • connect_timeout (float) –

    Timeout in seconds when connecting to a peer.

  • online_timeout (float | None) –

    Timeout in seconds to wait for the endpoint to connect to its home relay before logging a warning. If None, the endpoint does not wait (e.g., because relays are disabled).

from_config classmethod

from_config(config: EndpointP2PConfig) -> Self

Get the options for a peer-to-peer configuration.

The preset determines the discovery service, and the relay mode overrides the relays of the preset.

Source code in proxystore/endpoint/p2p/manager.py
@classmethod
def from_config(cls, config: EndpointP2PConfig) -> Self:
    """Get the options for a peer-to-peer configuration.

    The preset determines the discovery service, and the relay mode
    overrides the relays of the preset.
    """
    # The n0 preset uses n0's relays and discovery. The minimal preset
    # uses neither.
    if config.discovery == 'n0':
        preset = iroh.preset_n0()
        n0_relays = None
    else:
        preset = iroh.preset_minimal()
        n0_relays = iroh.RelayMode.default_mode()

    if config.relays == 'n0':
        return cls(preset=preset, relay_mode=n0_relays)
    if config.relays == 'none':
        # Without relays, there is no home relay to wait on.
        return cls(
            preset=preset,
            relay_mode=iroh.RelayMode.disabled(),
            online_timeout=None,
        )
    return cls(
        preset=preset,
        relay_mode=iroh.RelayMode.custom_from_urls(config.relays),
    )

CloseCode

Bases: IntEnum

Application error codes used when closing a peer connection.

SHUTDOWN class-attribute instance-attribute

SHUTDOWN = 0

Endpoint is shutting down.

NOT_ALLOWED class-attribute instance-attribute

NOT_ALLOWED = 1

Peer is not in the allowlist of the endpoint.

PeerManager

PeerManager(
    secret_key: SecretKey,
    policy: PeerPolicy,
    *,
    options: PeerOptions | None = None,
    max_request_size: int | None = None,
    addr_cache: PeerAddrCache | None = None,
)

Manager of connections to peer endpoints.

The manager binds an iroh endpoint using the secret key of the ProxyStore endpoint, accepts connections from peers, and sends requests to peers. Each request is sent on its own bidirectional stream. A connection to a peer is used in both directions: requests to the peer are sent on the most recent connection, whichever endpoint opened it, and requests from the peer are accepted on every connection. Two connections to a peer only exist if both endpoints connect to each other at the same time.

The manager only communicates with peers allowed by its PeerPolicy. Connections from other endpoints are refused, and requests to other endpoints fail with PeerNotAllowedError. The policy is checked on each connection and request, and connections to peers which are no longer allowed are closed.

Example
manager = PeerManager(secret_key, Allowlist(endpoint_dir.peers_path))
await manager.start(handler)
response = await manager.request(peer_id, Message(Op.GET, meta))
await manager.close()

Parameters:

  • secret_key (SecretKey) –

    Secret key of the endpoint.

  • policy (PeerPolicy) –

    Policy of which peers are allowed.

  • options (PeerOptions | None, default: None ) –

    Options of connections to peers. Defaults to PeerOptions().

  • max_request_size (int | None, default: None ) –

    Maximum size in bytes of the data in a request from a peer or None for no limit.

  • addr_cache (PeerAddrCache | None, default: None ) –

    Optional cache where the addresses of peers are saved (see PeerAddrCache). Cached addresses are used when connecting to peers so peers can be reached even if discovery is unavailable.

Source code in proxystore/endpoint/p2p/manager.py
def __init__(
    self,
    secret_key: SecretKey,
    policy: PeerPolicy,
    *,
    options: PeerOptions | None = None,
    max_request_size: int | None = None,
    addr_cache: PeerAddrCache | None = None,
) -> None:
    self._secret_key = secret_key
    self._id = secret_key.endpoint_id
    self._policy = policy
    self._options = PeerOptions() if options is None else options
    self._max_request_size = max_request_size
    self._addr_cache = addr_cache

    self._endpoint: iroh.Endpoint | None = None
    self._handler: RequestHandler | None = None
    self._accept_task: asyncio.Task[None] | None = None
    self._online_task: asyncio.Task[None] | None = None
    self._tasks: set[asyncio.Task[None]] = set()

    self._addr_hints: dict[EndpointId, iroh.EndpointAddr] = {}
    self._dial_locks: collections.defaultdict[
        EndpointId,
        asyncio.Lock,
    ] = collections.defaultdict(asyncio.Lock)
    # Open connections to each peer. Requests from the peer are
    # accepted on all of them.
    self._connections: dict[EndpointId, set[PeerConnection]] = {}
    # Connection used to send requests to each peer.
    self._preferred: dict[EndpointId, PeerConnection] = {}
    self._closed = False

id property

ID of this endpoint.

policy property

policy: PeerPolicy

Policy of which peers are allowed.

options property

options: PeerOptions

Options of connections to peers.

endpoint property

endpoint: Endpoint

Underlying iroh endpoint.

Raises:

path

path(peer_id: EndpointId) -> PathInfo | None

Get the path used by the connection to a peer.

Returns:

  • PathInfo | None –

    The path of the connection used to send requests to the peer or None if there is no open connection.

Source code in proxystore/endpoint/p2p/manager.py
def path(self, peer_id: EndpointId) -> PathInfo | None:
    """Get the path used by the connection to a peer.

    Returns:
        The path of the connection used to send requests to the peer or \
        `None` if there is no open connection.
    """
    connection = self._preferred.get(peer_id)
    if connection is None or connection.closed:
        return None
    return connection.path()

addr

addr() -> EndpointAddr

Get the current address of this endpoint.

The address contains the ID, home relay URL, and direct addresses of the endpoint. Other endpoints can use the address to connect to this endpoint without discovery (see add_peer_addr()).

Source code in proxystore/endpoint/p2p/manager.py
def addr(self) -> iroh.EndpointAddr:
    """Get the current address of this endpoint.

    The address contains the ID, home relay URL, and direct addresses
    of the endpoint. Other endpoints can use the address to connect to
    this endpoint without discovery (see
    [`add_peer_addr()`][proxystore.endpoint.p2p.manager.PeerManager.add_peer_addr]).
    """
    return self.endpoint.addr()

add_peer_addr

add_peer_addr(addr: EndpointAddr) -> None

Add a known address of a peer to use when connecting to the peer.

This is useful when discovery is unavailable.

Source code in proxystore/endpoint/p2p/manager.py
def add_peer_addr(self, addr: iroh.EndpointAddr) -> None:
    """Add a known address of a peer to use when connecting to the peer.

    This is useful when discovery is unavailable.
    """
    peer_id = EndpointId.from_str(str(addr.id()))
    self._addr_hints[peer_id] = addr

peer_name

peer_name(peer_id: EndpointId) -> str

Format the ID of a peer with its name for logging.

The name is from the peer policy or unknown if the peer has no name (see EndpointId.log_name()).

Source code in proxystore/endpoint/p2p/manager.py
def peer_name(self, peer_id: EndpointId) -> str:
    """Format the ID of a peer with its name for logging.

    The name is from the peer policy or `unknown` if the peer has no
    name (see
    [`EndpointId.log_name()`][proxystore.endpoint.identity.EndpointId.log_name]).
    """
    name = self._policy.name_of(peer_id)
    return peer_id.log_name('unknown' if name is None else name)

start async

start(handler: RequestHandler) -> None

Bind the endpoint and start accepting connections from peers.

Note

The iroh bindings run background threads so the manager must be started after the process is daemonized or forked.

Parameters:

Source code in proxystore/endpoint/p2p/manager.py
async def start(self, handler: RequestHandler) -> None:
    """Bind the endpoint and start accepting connections from peers.

    Note:
        The iroh bindings run background threads so the manager must be
        started after the process is daemonized or forked.

    Args:
        handler: Handler of requests from peers.
    """
    if self._endpoint is not None:
        return
    self._handler = handler
    if self._addr_cache is not None:
        cached = self._addr_cache.load()
        for peer_id, addr in cached.items():
            self._addr_hints.setdefault(peer_id, addr)
        logger.info(
            'Loaded %d cached peer address(es) from %s',
            len(cached),
            self._addr_cache.path,
        )
    # uniffi_set_event_loop() is intentionally not called. It sets a
    # process-wide event loop that the bindings then use for every call,
    # which breaks when a different event loop is used later. It is only
    # needed for callbacks from Rust into Python which are not used.
    options = self._options
    self._endpoint = await iroh.Endpoint.bind(
        iroh.EndpointOptions(
            preset=options.preset,
            secret_key=self._secret_key.to_bytes(),
            alpns=supported_alpns(),
            relay_mode=options.relay_mode,
            bind_addr=options.bind_addr,
        ),
    )
    self._accept_task = spawn_guarded_background_task(self._accept_loop)
    self._accept_task.set_name(f'peer-manager-{self.id}-accept')
    if options.online_timeout is not None:
        self._online_task = asyncio.create_task(
            self._wait_online(options.online_timeout),
        )
    logger.info(
        'Listening for peer connections on %s',
        ', '.join(self.endpoint.bound_sockets()),
    )

close async

close() -> None

Close all peer connections and the endpoint.

This is idempotent so it is safe to call multiple times.

Source code in proxystore/endpoint/p2p/manager.py
async def close(self) -> None:
    """Close all peer connections and the endpoint.

    This is idempotent so it is safe to call multiple times.
    """
    if self._closed:
        return
    self._closed = True
    # Connections are closed before the tasks are cancelled because
    # cancelling the task serving a connection forgets the connection
    # without closing it.
    for peer_id in list(self._connections):
        self._close_peer(peer_id, CloseCode.SHUTDOWN, b'shutdown')
    for task in (self._accept_task, self._online_task, *self._tasks):
        if task is not None:
            task.cancel()
            with contextlib.suppress(asyncio.CancelledError):
                await task
    if self._endpoint is not None:
        await self._endpoint.close()
    logger.info('Peer manager closed')

request async

request(peer_id: EndpointId, request: Message) -> Message

Send a request to a peer and wait for the response.

Parameters:

Returns:

  • Message –

    The response message.

Raises:

Source code in proxystore/endpoint/p2p/manager.py
async def request(
    self,
    peer_id: EndpointId,
    request: Message,
) -> Message:
    """Send a request to a peer and wait for the response.

    Args:
        peer_id: ID of the peer.
        request: Request message.

    Returns:
        The response message.

    Raises:
        PeerNotAllowedError: If the peer is not in the allowlist or the
            peer refused the connection.
        PeerConnectionTimeoutError: If connecting to the peer times out.
        PeerUnavailableError: If the request fails.
    """
    if not self._enforce_policy(peer_id):
        raise PeerNotAllowedError(
            f'Endpoint {peer_id} is not in the allowlist of peers. Add '
            'the endpoint with "proxystore-endpoint peers add".',
        )

    connection, fresh = await self._get_connection(peer_id)
    if not fresh:
        try:
            return await self._send(connection, request)
        except PeerUnavailableError as e:
            # A cached connection may have been closed (e.g., because
            # the peer restarted) so retry once with a new connection.
            logger.debug(
                'Retrying request to %s with a new connection: %s',
                self.peer_name(peer_id),
                e,
            )
        connection, _ = await self._get_connection(peer_id)
    return await self._send(connection, request)