Skip to content

proxystore.endpoint.peers

Peer endpoint allowlist.

An endpoint only communicates with the peer endpoints in its allowlist, the peers.toml file in the endpoint directory. Allowlisting is symmetric: two endpoints can only communicate if each endpoint lists the other.

peers.toml
[peers]
my-laptop = "5d3e...a1f2"
cluster = "bb04...9c3e"

The allowlist can be changed while the endpoint is running. The endpoint checks if the file changed at most once per second (see Allowlist) so a removed peer is denied access within about a second.

PEERS_VERSION module-attribute

PEERS_VERSION = 1

Format version of the peers file.

RELOAD_INTERVAL module-attribute

RELOAD_INTERVAL = 1.0

Default minimum seconds between checks for changes to the peers file.

PeersConfig

Bases: VersionedFile

Allowlist of peer endpoints.

Attributes:

  • version (int) –

    Format version of the peers file.

  • peers (dict[str, EndpointId]) –

    Mapping of peer names to endpoint IDs. Names are only used to help users manage the allowlist.

Raises:

  • ValueError –

    If a name does not contain only alphanumeric, dash, or underscore characters, if an endpoint ID is invalid, if an endpoint ID is listed more than once, or if the version is not supported.

name_of

name_of(endpoint_id: EndpointId) -> str | None

Get the name of the peer with the endpoint ID.

Source code in proxystore/endpoint/peers.py
def name_of(self, endpoint_id: EndpointId) -> str | None:
    """Get the name of the peer with the endpoint ID."""
    for name, peer_id in self.peers.items():
        if peer_id == endpoint_id:
            return name
    return None

Peers

Peers(path: str, *, owner_id: EndpointId | None = None)

Peers of an endpoint stored in its peers.toml file.

Example
peers = EndpointDir.from_name('my-ep').peers()
peers.add('laptop', 'ed92...')
assert peers.read().peers == {'laptop': 'ed92...'}
peers.remove('laptop')

Parameters:

  • path (str) –

    Path to the peers.toml file.

  • owner_id (EndpointId | None, default: None ) –

    ID of the endpoint which owns the peers. Used to prevent an endpoint from adding itself as a peer.

Source code in proxystore/endpoint/peers.py
def __init__(self, path: str, *, owner_id: EndpointId | None = None):
    self.path = path
    self.owner_id = owner_id

read

read() -> PeersConfig

Read the peers.

Returns:

  • PeersConfig –

    The peers or no peers if the file does not exist.

Raises:

Source code in proxystore/endpoint/peers.py
def read(self) -> PeersConfig:
    """Read the peers.

    Returns:
        The peers or no peers if the file does not exist.

    Raises:
        EndpointConfigError: If the file cannot be parsed or is invalid.
    """
    try:
        return read_model(PeersConfig, self.path)
    except FileNotFoundError:
        return PeersConfig()

write

write(peers: PeersConfig) -> None

Atomically write the peers.

Source code in proxystore/endpoint/peers.py
def write(self, peers: PeersConfig) -> None:
    """Atomically write the peers."""
    write_model(self.path, peers)

add

add(name: str, endpoint_id: str) -> EndpointId

Add a peer.

Parameters:

  • name (str) –

    Name of the peer.

  • endpoint_id (str) –

    ID of the peer endpoint.

Returns:

Raises:

  • PeerExistsError –

    If a peer with the name already exists.

  • EndpointConfigError –

    If the name or ID is invalid, the ID is the ID of the owner, the endpoint is already a peer with a different name, or the file cannot be parsed.

Source code in proxystore/endpoint/peers.py
def add(self, name: str, endpoint_id: str) -> EndpointId:
    """Add a peer.

    Args:
        name: Name of the peer.
        endpoint_id: ID of the peer endpoint.

    Returns:
        The ID of the peer.

    Raises:
        PeerExistsError: If a peer with the name already exists.
        EndpointConfigError: If the name or ID is invalid, the ID is the
            ID of the owner, the endpoint is already a peer with a
            different name, or the file cannot be parsed.
    """
    try:
        check_name(name, 'Peer')
        peer_id = EndpointId.from_str(endpoint_id)
    except ValueError as e:
        raise EndpointConfigError(str(e)) from None
    if peer_id == self.owner_id:
        raise EndpointConfigError(
            'An endpoint cannot be a peer of itself.',
        )
    peers = self.read()
    if name in peers.peers:
        raise PeerExistsError(f'A peer named {name} already exists.')
    existing = peers.name_of(peer_id)
    if existing is not None:
        raise EndpointConfigError(
            f'Endpoint {peer_id} is already a peer named {existing}.',
        )
    peers.peers[name] = peer_id
    self.write(peers)
    return peer_id

remove

remove(name: str) -> EndpointId

Remove a peer.

If the endpoint is running, the peer is denied access within about a second (see Allowlist).

Parameters:

  • name (str) –

    Name of the peer.

Returns:

Raises:

Source code in proxystore/endpoint/peers.py
def remove(self, name: str) -> EndpointId:
    """Remove a peer.

    If the endpoint is running, the peer is denied access within about
    a second (see [`Allowlist`][proxystore.endpoint.peers.Allowlist]).

    Args:
        name: Name of the peer.

    Returns:
        The ID of the removed peer.

    Raises:
        PeerNotFoundError: If there is no peer with the name.
        EndpointConfigError: If the file cannot be parsed.
    """
    peers = self.read()
    peer_id = peers.peers.pop(name, None)
    if peer_id is None:
        raise PeerNotFoundError(f'No peer named {name}.')
    self.write(peers)
    return peer_id

Allowlist

Allowlist(
    path: str, *, reload_interval: float = RELOAD_INTERVAL
)

Allowlist of peer endpoints backed by a peers.toml file.

This is the PeerPolicy of endpoints. The file is reloaded when it changes. A missing file is an empty allowlist. If the file is malformed, all peers are denied until it is fixed.

Checking if the file changed requires a stat() call which can be slow on network file systems, so the file is checked at most once every reload_interval seconds rather than on every request.

Parameters:

  • path (str) –

    Path to the peers.toml file.

  • reload_interval (float, default: RELOAD_INTERVAL ) –

    Minimum seconds between checks for changes to the file. If 0, the file is checked each time the allowlist is used.

Source code in proxystore/endpoint/peers.py
def __init__(
    self,
    path: str,
    *,
    reload_interval: float = RELOAD_INTERVAL,
) -> None:
    self.path = path
    self.reload_interval = reload_interval
    self._checked: float | None = None
    self._state: _FileState | None = None
    self._peers = PeersConfig()

peers property

peers: PeersConfig

Current allowlist, reloaded if the file changed.

reload

reload(*, force: bool = False) -> None

Reload the allowlist if the file changed.

Parameters:

  • force (bool, default: False ) –

    Check if the file changed even if the reload interval has not passed since the last check.

Source code in proxystore/endpoint/peers.py
def reload(self, *, force: bool = False) -> None:
    """Reload the allowlist if the file changed.

    Args:
        force: Check if the file changed even if the reload interval has
            not passed since the last check.
    """
    now = time.monotonic()
    if (
        not force
        and self._checked is not None
        and now - self._checked < self.reload_interval
    ):
        return
    self._checked = now

    try:
        stat = os.stat(self.path)
    except FileNotFoundError:
        state = None
    else:
        state = _FileState(stat.st_mtime_ns, stat.st_size, stat.st_ino)

    if state == self._state:
        return
    self._state = state

    if state is None:
        self._peers = PeersConfig()
    else:
        try:
            self._peers = Peers(self.path).read()
        except ValueError as e:
            logger.error(
                'Failed to load peer allowlist from %s. All peers will '
                'be denied until the file is fixed: %s',
                self.path,
                e,
            )
            self._peers = PeersConfig()

allowed

allowed(endpoint_id: EndpointId) -> bool

Check if the endpoint is in the allowlist.

Source code in proxystore/endpoint/peers.py
def allowed(self, endpoint_id: EndpointId) -> bool:
    """Check if the endpoint is in the allowlist."""
    return endpoint_id in self.peers.peers.values()

name_of

name_of(endpoint_id: EndpointId) -> str | None

Get the name of the peer or None if it is not allowed.

Source code in proxystore/endpoint/peers.py
def name_of(self, endpoint_id: EndpointId) -> str | None:
    """Get the name of the peer or `None` if it is not allowed."""
    return self.peers.name_of(endpoint_id)