Skip to content

proxystore.endpoint.client

Client for communicating with a local endpoint.

Note

Clients communicate with endpoints on the local network over TCP using the protocol defined in proxystore.endpoint.protocol. It is not intended that clients from outside the local network interact with an endpoint this way. (Rather, they should connect to their own local endpoint, which peers with remote endpoints.)

REQUEST_TIMEOUT module-attribute

REQUEST_TIMEOUT = 60

Default seconds a request to the local endpoint can go without progress.

EndpointClient

EndpointClient(
    sock: socket,
    info: EndpointInfo,
    *,
    protocol_version: int,
    request_timeout: float | None = REQUEST_TIMEOUT,
)

Connection to a local endpoint.

Use from_name() to connect to a local endpoint by name, or connect() to connect to an address directly.

Warning

A client is not thread-safe because a connection can only process one request at a time. Use a separate client per thread.

Example
with EndpointClient.from_name('my-endpoint') as client:
    client.set('key', b'value')
    assert client.get('key') == b'value'

Parameters:

  • sock (socket) –

    Connected socket that has completed the handshake.

  • info (EndpointInfo) –

    Information about the endpoint.

  • protocol_version (int) –

    Protocol version negotiated in the handshake.

  • request_timeout (float | None, default: REQUEST_TIMEOUT ) –

    Seconds a request handled by this endpoint can go without sending or receiving any data, including while the endpoint handles the request (e.g., writes the object to its storage), before it fails with an EndpointTimeoutError. Requests forwarded to a peer endpoint have no timeout in the client because the endpoint does not respond until the peer does. If None, requests have no timeout.

Raises:

  • ValueError –

    If request_timeout is not positive.

Source code in proxystore/endpoint/client.py
def __init__(
    self,
    sock: socket.socket,
    info: EndpointInfo,
    *,
    protocol_version: int,
    request_timeout: float | None = REQUEST_TIMEOUT,
) -> None:
    self._socket = sock
    self.info = info
    self.protocol_version = protocol_version
    self.request_timeout = check_request_timeout(request_timeout)
    self.closed = False
    self._next_request_id = 1

connect classmethod

connect(
    host: str,
    port: int,
    token: EndpointToken,
    *,
    tls_fingerprint: str | None = None,
    timeout: float | None = 10,
    request_timeout: float | None = REQUEST_TIMEOUT,
) -> Self

Connect to an endpoint and complete the handshake.

Parameters:

  • host (str) –

    Host address of the endpoint.

  • port (int) –

    Port of the endpoint.

  • token (EndpointToken) –

    Token of the endpoint (see EndpointDir.read_connection()).

  • tls_fingerprint (str | None, default: None ) –

    SHA-256 fingerprint of the endpoint's TLS certificate (see EndpointDir.read_connection()). If provided, the connection is encrypted with TLS and the endpoint's certificate must match the fingerprint.

  • timeout (float | None, default: 10 ) –

    Timeout in seconds for connecting and completing the handshake.

  • request_timeout (float | None, default: REQUEST_TIMEOUT ) –

    Seconds a request can go without progress (see EndpointClient).

Warns:

  • VersionMismatchWarning –

    If the endpoint uses a different ProxyStore version or Python minor version than this client.

Raises:

Source code in proxystore/endpoint/client.py
@classmethod
def connect(
    cls,
    host: str,
    port: int,
    token: EndpointToken,
    *,
    tls_fingerprint: str | None = None,
    timeout: float | None = 10,
    request_timeout: float | None = REQUEST_TIMEOUT,
) -> Self:
    """Connect to an endpoint and complete the handshake.

    Args:
        host: Host address of the endpoint.
        port: Port of the endpoint.
        token: Token of the endpoint (see
            [`EndpointDir.read_connection()`][proxystore.endpoint.directory.EndpointDir.read_connection]).
        tls_fingerprint: SHA-256 fingerprint of the endpoint's TLS
            certificate (see
            [`EndpointDir.read_connection()`][proxystore.endpoint.directory.EndpointDir.read_connection]).
            If provided, the connection is encrypted with TLS and the
            endpoint's certificate must match the fingerprint.
        timeout: Timeout in seconds for connecting and completing the
            handshake.
        request_timeout: Seconds a request can go without progress (see
            [`EndpointClient`][proxystore.endpoint.client.EndpointClient]).

    Warns:
        VersionMismatchWarning: If the endpoint uses a different
            ProxyStore version or Python minor version than this client.

    Raises:
        ValueError: If `request_timeout` is not positive.
        EndpointNotRunningError: If the connection is refused.
        EndpointConnectionError: If the connection cannot be established
            or is lost during the handshake (e.g., a timeout).
        EndpointAuthError: If the client or endpoint fails
            authentication.
        EndpointProtocolError: If the endpoint uses an incompatible
            protocol.
    """
    check_request_timeout(request_timeout)
    try:
        sock = socket.create_connection((host, port), timeout=timeout)
    except ConnectionRefusedError as e:
        raise EndpointNotRunningError(
            f'Connection to the endpoint at {host}:{port} was refused. '
            'Is the endpoint running?',
        ) from e
    except OSError as e:
        raise EndpointConnectionError(
            f'Unable to connect to the endpoint at {host}:{port}: {e}',
        ) from e

    try:
        sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
        _enable_keepalive(sock)
        if tls_fingerprint is not None:
            sock = _wrap_tls(sock, tls_fingerprint)
        info, version = _handshake(sock, token)
        sock.settimeout(None)
    except OSError as e:
        sock.close()
        raise EndpointConnectionError(
            f'Lost connection to the endpoint at {host}:{port} during '
            f'the handshake: {e}',
        ) from e
    except BaseException:
        sock.close()
        raise

    warning = Versions.current().mismatch_warning(info.versions)
    if warning is not None:
        warnings.warn(
            f'Endpoint {info.name} ({info.id.short()}) uses different '
            f'versions than this client: {warning}',
            VersionMismatchWarning,
            stacklevel=2,
        )
    logger.debug(
        'Connected to endpoint %s at %s:%s (tls=%s, proxystore=%s, '
        'python=%s)',
        info.id.log_name(info.name),
        host,
        port,
        tls_fingerprint is not None,
        info.versions.proxystore,
        info.versions.python,
    )
    return cls(
        sock,
        info,
        protocol_version=version,
        request_timeout=request_timeout,
    )

from_dir classmethod

from_dir(
    endpoint_dir: EndpointDir,
    *,
    timeout: float | None = 10,
    request_timeout: float | None = REQUEST_TIMEOUT,
) -> Self

Connect to a local endpoint using its connection file.

The connection file is read each time because the endpoint writes a new one each time it starts.

Parameters:

  • endpoint_dir (EndpointDir) –

    Directory of the endpoint.

  • timeout (float | None, default: 10 ) –

    Timeout in seconds for connecting and completing the handshake.

  • request_timeout (float | None, default: REQUEST_TIMEOUT ) –

    Seconds a request can go without progress (see EndpointClient).

Raises:

  • EndpointNotFoundError –

    If the endpoint directory does not exist.

  • EndpointNotRunningError –

    If the endpoint's connection file does not exist (i.e., the endpoint is not running).

  • EndpointAuthError –

    If the connection file cannot be read or is malformed.

  • EndpointError –

    If the connection or handshake fails (see connect()).

Source code in proxystore/endpoint/client.py
@classmethod
def from_dir(
    cls,
    endpoint_dir: EndpointDir,
    *,
    timeout: float | None = 10,
    request_timeout: float | None = REQUEST_TIMEOUT,
) -> Self:
    """Connect to a local endpoint using its connection file.

    The connection file is read each time because the endpoint writes
    a new one each time it starts.

    Args:
        endpoint_dir: Directory of the endpoint.
        timeout: Timeout in seconds for connecting and completing the
            handshake.
        request_timeout: Seconds a request can go without progress (see
            [`EndpointClient`][proxystore.endpoint.client.EndpointClient]).

    Raises:
        EndpointNotFoundError: If the endpoint directory does not exist.
        EndpointNotRunningError: If the endpoint's connection file does
            not exist (i.e., the endpoint is not running).
        EndpointAuthError: If the connection file cannot be read or is
            malformed.
        EndpointError: If the connection or handshake fails (see
            [`connect()`][proxystore.endpoint.client.EndpointClient.connect]).
    """
    endpoint_dir.check_exists()
    try:
        info = endpoint_dir.read_connection()
    except FileNotFoundError as e:
        raise EndpointNotRunningError(
            _missing_connection_file_message(endpoint_dir),
        ) from e
    except (OSError, ValueError) as e:
        raise EndpointAuthError(
            f'Unable to read the connection file of the endpoint in '
            f'{endpoint_dir}: {e}',
        ) from e
    return cls.connect(
        info.host,
        info.port,
        info.token,
        tls_fingerprint=info.tls_fingerprint,
        timeout=timeout,
        request_timeout=request_timeout,
    )

from_name classmethod

from_name(
    name: str,
    *,
    proxystore_dir: str | None = None,
    timeout: float | None = 10,
    request_timeout: float | None = REQUEST_TIMEOUT,
) -> Self

Connect to a local endpoint by name.

Parameters:

  • name (str) –

    Name of the endpoint.

  • proxystore_dir (str | None, default: None ) –

    ProxyStore home directory containing the endpoint. Defaults to home_dir().

  • timeout (float | None, default: 10 ) –

    Timeout in seconds for connecting and completing the handshake.

  • request_timeout (float | None, default: REQUEST_TIMEOUT ) –

    Seconds a request can go without progress (see EndpointClient).

Raises:

  • EndpointNotFoundError –

    If no endpoint with the name exists.

  • EndpointError –

    If connecting to the endpoint fails (see from_dir()).

Source code in proxystore/endpoint/client.py
@classmethod
def from_name(
    cls,
    name: str,
    *,
    proxystore_dir: str | None = None,
    timeout: float | None = 10,
    request_timeout: float | None = REQUEST_TIMEOUT,
) -> Self:
    """Connect to a local endpoint by name.

    Args:
        name: Name of the endpoint.
        proxystore_dir: ProxyStore home directory containing the
            endpoint. Defaults to
            [`home_dir()`][proxystore.utils.environment.home_dir].
        timeout: Timeout in seconds for connecting and completing the
            handshake.
        request_timeout: Seconds a request can go without progress (see
            [`EndpointClient`][proxystore.endpoint.client.EndpointClient]).

    Raises:
        EndpointNotFoundError: If no endpoint with the name exists.
        EndpointError: If connecting to the endpoint fails (see
            [`from_dir()`][proxystore.endpoint.client.EndpointClient.from_dir]).
    """
    endpoint_dir = EndpointDir.from_name(name, proxystore_dir)
    return cls.from_dir(
        endpoint_dir,
        timeout=timeout,
        request_timeout=request_timeout,
    )

close

close() -> None

Close the connection.

Source code in proxystore/endpoint/client.py
def close(self) -> None:
    """Close the connection."""
    if not self.closed:
        self.closed = True
        self._socket.close()
        logger.debug(
            'Closed connection to endpoint %s',
            self.info.id.log_name(self.info.name),
        )

evict

evict(key: str, target: str | None = None) -> None

Evict the object associated with the key.

Parameters:

  • key (str) –

    Key associated with object to evict.

  • target (str | None, default: None ) –

    Optional ID of a peer endpoint to forward the operation to.

Raises:

Source code in proxystore/endpoint/client.py
def evict(self, key: str, target: str | None = None) -> None:
    """Evict the object associated with the key.

    Args:
        key: Key associated with object to evict.
        target: Optional ID of a peer endpoint to forward the operation to.

    Raises:
        ValueError: If `target` is not a valid endpoint ID.
        EndpointError: If the request fails.
    """
    self._request(Op.EVICT, _request(key, target))

exists

exists(key: str, target: str | None = None) -> bool

Check if an object associated with the key exists.

Parameters:

  • key (str) –

    Key potentially associated with stored object.

  • target (str | None, default: None ) –

    Optional ID of a peer endpoint to forward the operation to.

Returns:

  • bool –

    If an object associated with the key exists.

Raises:

Source code in proxystore/endpoint/client.py
def exists(self, key: str, target: str | None = None) -> bool:
    """Check if an object associated with the key exists.

    Args:
        key: Key potentially associated with stored object.
        target: Optional ID of a peer endpoint to forward the operation to.

    Returns:
        If an object associated with the key exists.

    Raises:
        ValueError: If `target` is not a valid endpoint ID.
        EndpointError: If the request fails.
    """
    response = self._request(Op.EXISTS, _request(key, target))
    return ExistsResult.decode(response.meta).exists

get

get(
    key: str, target: str | None = None
) -> bytearray | None

Get the serialized object associated with the key.

Parameters:

  • key (str) –

    Key associated with object to retrieve.

  • target (str | None, default: None ) –

    Optional ID of a peer endpoint to forward the operation to.

Returns:

  • bytearray | None –

    Serialized object or None if the object does not exist.

Raises:

Source code in proxystore/endpoint/client.py
def get(
    self,
    key: str,
    target: str | None = None,
) -> bytearray | None:
    """Get the serialized object associated with the key.

    Args:
        key: Key associated with object to retrieve.
        target: Optional ID of a peer endpoint to forward the operation to.

    Returns:
        Serialized object or `None` if the object does not exist.

    Raises:
        ValueError: If `target` is not a valid endpoint ID.
        EndpointError: If the request fails.
    """
    response = self._request(Op.GET, _request(key, target))
    if response.code == Status.NOT_FOUND:
        return None
    data = response.data
    # Data is always read into a bytearray unless it is empty.
    return data if isinstance(data, bytearray) else bytearray(data)

set

set(
    key: str, data: BytesLike, target: str | None = None
) -> None

Set the serialized object associated with the key.

Parameters:

  • key (str) –

    Key to associate with the object.

  • data (BytesLike) –

    Serialized object.

  • target (str | None, default: None ) –

    Optional ID of a peer endpoint to forward the operation to.

Raises:

Source code in proxystore/endpoint/client.py
def set(
    self,
    key: str,
    data: BytesLike,
    target: str | None = None,
) -> None:
    """Set the serialized object associated with the key.

    Args:
        key: Key to associate with the object.
        data: Serialized object.
        target: Optional ID of a peer endpoint to forward the operation to.

    Raises:
        ObjectSizeExceededError: If the size of `data` exceeds the
            maximum object size of the endpoint.
        ValueError: If `target` is not a valid endpoint ID.
        EndpointError: If the request fails.
    """
    size = memoryview(data).nbytes
    max_size = self.info.max_object_size
    if max_size is not None and size > max_size:
        raise ObjectSizeExceededError(
            f'Data size ({size} bytes) exceeds the maximum object size '
            f'of the endpoint ({max_size} bytes).',
        )
    self._request(Op.SET, _request(key, target), data)

ping

ping(target: str | None = None) -> PingResult

Measure the latency of and path to a peer endpoint.

The local endpoint sends a request to the peer and reports the time until it received the response and the network path of the connection. The first ping to a peer includes the time to establish the connection.

Parameters:

  • target (str | None, default: None ) –

    Optional ID of the peer endpoint to ping. If None, the local endpoint is pinged.

Raises:

Source code in proxystore/endpoint/client.py
def ping(self, target: str | None = None) -> PingResult:
    """Measure the latency of and path to a peer endpoint.

    The local endpoint sends a request to the peer and reports the time
    until it received the response and the network path of the
    connection. The first ping to a peer includes the time to establish
    the connection.

    Args:
        target: Optional ID of the peer endpoint to ping. If `None`,
            the local endpoint is pinged.

    Raises:
        ValueError: If `target` is not a valid endpoint ID.
        EndpointError: If the request fails.
    """
    response = self._request(Op.PING, _request(None, target))
    return PingResult.decode(response.meta)

check_request_timeout

check_request_timeout(
    timeout: float | None,
) -> float | None

Check that a request timeout is positive or None.

Raises:

Source code in proxystore/endpoint/client.py
def check_request_timeout(timeout: float | None) -> float | None:
    """Check that a request timeout is positive or `None`.

    Raises:
        ValueError: If `timeout` is not positive.
    """
    # A timeout of 0 would make the socket non-blocking.
    if timeout is not None and timeout <= 0:
        raise ValueError(
            f'The request timeout must be positive or None. Got {timeout}.',
        )
    return timeout