Skip to content

proxystore.endpoint.process

Run, start, and stop endpoint processes.

serve() runs an Endpoint in the current process until it receives a signal. start_endpoint() runs an endpoint in this process or as a daemon, and stop_endpoint() stops the process of a running endpoint using the PID in its connection file. These functions raise EndpointError subclasses on failure.

Note

This module requires the endpoints extra.

start_endpoint

start_endpoint(
    endpoint_dir: EndpointDir,
    *,
    detach: bool = False,
    log_level: int | str = INFO,
) -> None

Start an endpoint.

Warning

If detach is False, this function does not return until the endpoint receives SIGINT or SIGTERM.

Parameters:

  • endpoint_dir (EndpointDir) –

    Directory of the endpoint to start.

  • detach (bool, default: False ) –

    Start the endpoint as a daemon process.

  • log_level (int | str, default: INFO ) –

    Logging level of the endpoint.

Raises:

  • EndpointNotFoundError –

    If the endpoint does not exist.

  • EndpointConfigError –

    If the configuration is invalid or the host address cannot be resolved.

  • EndpointRunningError –

    If the endpoint is already running on this or another host.

Source code in proxystore/endpoint/process.py
def start_endpoint(
    endpoint_dir: EndpointDir,
    *,
    detach: bool = False,
    log_level: int | str = logging.INFO,
) -> None:
    """Start an endpoint.

    Warning:
        If `detach` is `False`, this function does not return until the
        endpoint receives SIGINT or SIGTERM.

    Args:
        endpoint_dir: Directory of the endpoint to start.
        detach: Start the endpoint as a daemon process.
        log_level: Logging level of the endpoint.

    Raises:
        EndpointNotFoundError: If the endpoint does not exist.
        EndpointConfigError: If the configuration is invalid or the host
            address cannot be resolved.
        EndpointRunningError: If the endpoint is already running on this or
            another host.
    """
    # These checks are repeated by the endpoint when it starts but are
    # checked first so errors are raised to the caller rather than only
    # written to the log of a daemon.
    endpoint_dir.check_stopped()
    host = endpoint_dir.read_config().host
    try:
        resolve_host(host)
    except OSError as e:
        raise EndpointConfigError(
            f'Unable to resolve the host address ({host}): {e}',
        ) from e

    context: contextlib.AbstractContextManager[object]
    if detach:
        logger.info('Starting endpoint process as daemon')
        logger.info('Logs will be written to %s', endpoint_dir.log_path)
        context = daemon.DaemonContext(
            working_directory=endpoint_dir.path,
            umask=0o077,
            detach_process=True,
            # Note: stdin, stdout, stderr left as None which binds to /dev/null
        )
    else:
        context = contextlib.nullcontext()

    with context:
        # Logging is configured after daemonizing because the daemon closes
        # all open files (e.g., the log file).
        configure_logging(log_level, endpoint_dir.log_path)
        # Note: serve will handle most interrupts which can be reasonably
        # handled and return gracefully.
        serve(endpoint_dir)

configure_logging

configure_logging(
    log_level: int | str, log_file: str
) -> None

Configure logging of an endpoint process.

This sets the level of the root logger and appends the log to a file. Existing handlers of the root logger (e.g., of the CLI) are unchanged. This is only called by start_endpoint() because it changes the logging of the entire process.

Parameters:

  • log_level (int | str) –

    Logging level.

  • log_file (str) –

    File path to append the log to. The parent directory is created if it does not exist.

Source code in proxystore/endpoint/process.py
def configure_logging(log_level: int | str, log_file: str) -> None:
    """Configure logging of an endpoint process.

    This sets the level of the root logger and appends the log to a file.
    Existing handlers of the root logger (e.g., of the CLI) are unchanged.
    This is only called by
    [`start_endpoint()`][proxystore.endpoint.process.start_endpoint] because
    it changes the logging of the entire process.

    Args:
        log_level: Logging level.
        log_file: File path to append the log to. The parent directory is
            created if it does not exist.
    """
    os.makedirs(os.path.dirname(log_file) or '.', exist_ok=True)
    handler = logging.FileHandler(log_file)
    handler.setFormatter(
        logging.Formatter(
            '[%(asctime)s.%(msecs)03d] %(levelname)-5s (%(name)s) :: '
            '%(message)s',
            datefmt='%Y-%m-%d %H:%M:%S',
        ),
    )
    root = logging.getLogger()
    root.addHandler(handler)
    root.setLevel(log_level)

serve

serve(
    endpoint_dir: EndpointDir, *, use_uvloop: bool = True
) -> None

Run an endpoint in the current process.

Warning

This function does not return until the process receives SIGINT or SIGTERM.

Parameters:

  • endpoint_dir (EndpointDir) –

    Directory of the endpoint with its configuration. The connection file is written to this directory while the endpoint is running.

  • use_uvloop (bool, default: True ) –

    Use uvloop as the event loop implementation.

Source code in proxystore/endpoint/process.py
def serve(endpoint_dir: EndpointDir, *, use_uvloop: bool = True) -> None:
    """Run an endpoint in the current process.

    Warning:
        This function does not return until the process receives SIGINT or
        SIGTERM.

    Args:
        endpoint_dir: Directory of the endpoint with its configuration. The
            connection file is written to this directory while the
            endpoint is running.
        use_uvloop: Use uvloop as the event loop implementation.
    """
    try:
        if use_uvloop:  # pragma: no cover
            logger.info('Using uvloop as the event loop')
            uvloop.run(_serve_async(endpoint_dir))
        else:
            asyncio.run(_serve_async(endpoint_dir))
    except Exception:
        # Intercept exception so we can log it in the case that the endpoint
        # is running as a daemon process. Otherwise the user will never see
        # the exception.
        logger.exception('Endpoint failed with an unhandled exception')
        raise
    except KeyboardInterrupt:  # pragma: no cover
        # SIGINT is handled by _serve_async once the event loop is running,
        # but can still be raised before then.
        pass
    finally:
        logger.info('Finished serving endpoint in %s', endpoint_dir)

stop_endpoint

stop_endpoint(
    endpoint_dir: EndpointDir, *, timeout: float = 5
) -> bool

Stop an endpoint running on this host.

The endpoint is sent SIGTERM and killed if it does not exit within timeout seconds.

Parameters:

  • endpoint_dir (EndpointDir) –

    Directory of the endpoint to stop.

  • timeout (float, default: 5 ) –

    Seconds to wait for the endpoint to exit.

Returns:

  • bool –

    True if the endpoint was running and was stopped, or False if the endpoint was not running.

Raises:

  • EndpointNotFoundError –

    If the endpoint does not exist.

  • EndpointRunningError –

    If the endpoint may be running on another host or is still starting.

Source code in proxystore/endpoint/process.py
def stop_endpoint(endpoint_dir: EndpointDir, *, timeout: float = 5) -> bool:
    """Stop an endpoint running on this host.

    The endpoint is sent SIGTERM and killed if it does not exit within
    `timeout` seconds.

    Args:
        endpoint_dir: Directory of the endpoint to stop.
        timeout: Seconds to wait for the endpoint to exit.

    Returns:
        `True` if the endpoint was running and was stopped, or `False` if \
        the endpoint was not running.

    Raises:
        EndpointNotFoundError: If the endpoint does not exist.
        EndpointRunningError: If the endpoint may be running on another host
            or is still starting.
    """
    if endpoint_dir.status() != EndpointStatus.RUNNING:
        # Raises if the endpoint may be running on another host and removes
        # a stale connection file.
        endpoint_dir.check_stopped()
        return False

    try:
        info = endpoint_dir.read_connection()
    except FileNotFoundError:
        raise EndpointRunningError(
            f'Endpoint {endpoint_dir.name} is running but has not written '
            'its connection file so it is likely still starting. Try again '
            'once it has started.',
        ) from None

    logger.debug('Terminating endpoint process (PID: %s)', info.pid)
    with contextlib.suppress(ProcessLookupError):
        os.kill(info.pid, signal.SIGTERM)

    if not _wait_for_exit(info.pid, timeout=timeout):  # pragma: no cover
        logger.warning(
            'Killing endpoint process (PID: %s) which did not exit within '
            '%s seconds',
            info.pid,
            timeout,
        )
        with contextlib.suppress(ProcessLookupError):
            os.kill(info.pid, signal.SIGKILL)

    # The endpoint removes its connection file when it stops unless it was
    # killed.
    endpoint_dir.remove_connection(info)
    return True