Skip to content

guidellm.scheduler.environments

Environment abstractions for coordinating scheduler execution across distributed nodes.

Provides abstractions that handle synchronization, timing coordination, error propagation, and lifecycle management for scheduler execution across single or multiple nodes. The Environment protocol defines the interface for distributed coordination while NonDistributedEnvironment provides a minimal implementation for single-node execution. Environments manage the complete execution lifecycle from parameter distribution through result aggregation.

Execution Flow: 1. sync_run_params() - Distribute workload and synchronize parameters 2. sync_run_start() - Coordinate synchronized start time 3. update_run_iteration() - Update state after each request iteration 4. sync_run_error() - Handle and propagate errors across nodes 5. sync_run_end() - Aggregate results and finalize execution

Environment

Bases: ABC, Generic[RequestT, ResponseT], InfoMixin

Abstract interface for coordinating scheduler execution across distributed nodes.

Defines the protocol for managing distributed scheduler execution including parameter synchronization, timing coordination, state updates, error propagation, and result aggregation. Implementations handle distributed coordination complexity while providing a unified interface for scheduler orchestration.

Source code in src/guidellm/scheduler/environments.py
class Environment(ABC, Generic[RequestT, ResponseT], InfoMixin):
    """
    Abstract interface for coordinating scheduler execution across distributed nodes.

    Defines the protocol for managing distributed scheduler execution including
    parameter synchronization, timing coordination, state updates, error propagation,
    and result aggregation. Implementations handle distributed coordination complexity
    while providing a unified interface for scheduler orchestration.
    """

    @abstractmethod
    async def sync_run_params(
        self,
        requests: DatasetIterT[RequestT],
        strategy: SchedulingStrategy,
        constraints: dict[str, Constraint | ConstraintInitializer],
    ) -> tuple[
        DatasetIterT[RequestT],
        SchedulingStrategy,
        dict[str, Constraint | ConstraintInitializer],
    ]:
        """
        Synchronize execution parameters across nodes and resolve local scope.

        :param requests: Complete set of requests to process across all nodes
        :param strategy: Scheduling strategy to apply during execution
        :param constraints: Runtime constraints to enforce during execution
        :return: Tuple of (local_requests, strategy, constraints) for this node
        :raises Exception: If parameter synchronization fails or nodes inconsistent
        """
        ...

    @abstractmethod
    async def sync_run_start(self) -> float:
        """
        Coordinate synchronized start time across all nodes.

        :return: Unix timestamp when all nodes should begin processing
        :raises Exception: If startup synchronization fails across nodes
        """
        ...

    @abstractmethod
    async def update_run_iteration(
        self,
        response: ResponseT | None,
        request: RequestT,
        request_info: RequestInfo,
        state: SchedulerState,
    ):
        """
        Update environment state with completed request iteration results.

        :param response: Response generated for the request, if successful
        :param request: The processed request
        :param request_info: Metadata about request processing including timings
        :param state: Current scheduler state with metrics and progress
        :raises Exception: If state update fails or indicates critical errors
        """
        ...

    @abstractmethod
    async def sync_run_error(self, err: list[Exception] | Exception):
        """
        Handle and propagate errors across all active nodes.

        :param err: The exception(s) that occurred during execution
        """
        ...

    @abstractmethod
    async def sync_run_end(
        self,
    ) -> AsyncIterator[
        tuple[
            ResponseT | None,
            RequestT,
            RequestInfo,
            SchedulerState,
        ]
    ]:
        """
        Finalize execution and aggregate results from all nodes.

        :return: Iterator of (response, request, request_info, state) tuples from
            remote nodes in distributed environments, empty for non-distributed
        :raises Exception: Any errors that occurred during execution
        """
        yield None  # type: ignore[misc]

sync_run_end() abstractmethod async

Finalize execution and aggregate results from all nodes.

Returns:

Type Description
AsyncIterator[tuple[ResponseT | None, RequestT, RequestInfo, SchedulerState]]

Iterator of (response, request, request_info, state) tuples from remote nodes in distributed environments, empty for non-distributed

Raises:

Type Description
Exception

Any errors that occurred during execution

Source code in src/guidellm/scheduler/environments.py
@abstractmethod
async def sync_run_end(
    self,
) -> AsyncIterator[
    tuple[
        ResponseT | None,
        RequestT,
        RequestInfo,
        SchedulerState,
    ]
]:
    """
    Finalize execution and aggregate results from all nodes.

    :return: Iterator of (response, request, request_info, state) tuples from
        remote nodes in distributed environments, empty for non-distributed
    :raises Exception: Any errors that occurred during execution
    """
    yield None  # type: ignore[misc]

sync_run_error(err) abstractmethod async

Handle and propagate errors across all active nodes.

Parameters:

Name Type Description Default
err list[Exception] | Exception

The exception(s) that occurred during execution

required
Source code in src/guidellm/scheduler/environments.py
@abstractmethod
async def sync_run_error(self, err: list[Exception] | Exception):
    """
    Handle and propagate errors across all active nodes.

    :param err: The exception(s) that occurred during execution
    """
    ...

sync_run_params(requests, strategy, constraints) abstractmethod async

Synchronize execution parameters across nodes and resolve local scope.

Parameters:

Name Type Description Default
requests DatasetIterT[RequestT]

Complete set of requests to process across all nodes

required
strategy SchedulingStrategy

Scheduling strategy to apply during execution

required
constraints dict[str, Constraint | ConstraintInitializer]

Runtime constraints to enforce during execution

required

Returns:

Type Description
tuple[DatasetIterT[RequestT], SchedulingStrategy, dict[str, Constraint | ConstraintInitializer]]

Tuple of (local_requests, strategy, constraints) for this node

Raises:

Type Description
Exception

If parameter synchronization fails or nodes inconsistent

Source code in src/guidellm/scheduler/environments.py
@abstractmethod
async def sync_run_params(
    self,
    requests: DatasetIterT[RequestT],
    strategy: SchedulingStrategy,
    constraints: dict[str, Constraint | ConstraintInitializer],
) -> tuple[
    DatasetIterT[RequestT],
    SchedulingStrategy,
    dict[str, Constraint | ConstraintInitializer],
]:
    """
    Synchronize execution parameters across nodes and resolve local scope.

    :param requests: Complete set of requests to process across all nodes
    :param strategy: Scheduling strategy to apply during execution
    :param constraints: Runtime constraints to enforce during execution
    :return: Tuple of (local_requests, strategy, constraints) for this node
    :raises Exception: If parameter synchronization fails or nodes inconsistent
    """
    ...

sync_run_start() abstractmethod async

Coordinate synchronized start time across all nodes.

Returns:

Type Description
float

Unix timestamp when all nodes should begin processing

Raises:

Type Description
Exception

If startup synchronization fails across nodes

Source code in src/guidellm/scheduler/environments.py
@abstractmethod
async def sync_run_start(self) -> float:
    """
    Coordinate synchronized start time across all nodes.

    :return: Unix timestamp when all nodes should begin processing
    :raises Exception: If startup synchronization fails across nodes
    """
    ...

update_run_iteration(response, request, request_info, state) abstractmethod async

Update environment state with completed request iteration results.

Parameters:

Name Type Description Default
response ResponseT | None

Response generated for the request, if successful

required
request RequestT

The processed request

required
request_info RequestInfo

Metadata about request processing including timings

required
state SchedulerState

Current scheduler state with metrics and progress

required

Raises:

Type Description
Exception

If state update fails or indicates critical errors

Source code in src/guidellm/scheduler/environments.py
@abstractmethod
async def update_run_iteration(
    self,
    response: ResponseT | None,
    request: RequestT,
    request_info: RequestInfo,
    state: SchedulerState,
):
    """
    Update environment state with completed request iteration results.

    :param response: Response generated for the request, if successful
    :param request: The processed request
    :param request_info: Metadata about request processing including timings
    :param state: Current scheduler state with metrics and progress
    :raises Exception: If state update fails or indicates critical errors
    """
    ...

NonDistributedEnvironment

Bases: Environment[RequestT, ResponseT]

Single-node scheduler execution environment with minimal coordination overhead.

Implements the Environment interface with no-op synchronization for local testing, development, and single-machine benchmarking. All synchronization methods return immediately without distributed coordination logic.

Example: :: from guidellm.scheduler import ( MaxNumberConstraint, MaxRequestsConstraintArgs, NonDistributedEnvironment, RequestInfo, SchedulerState, SynchronousStrategy, )

env = NonDistributedEnvironment()
requests = [f"req_{ind}" for ind in range(5)]
strategy = SynchronousStrategy()
args = MaxRequestsConstraintArgs(max_num=5)
constraints = {"max_requests": MaxNumberConstraint(args=args)}
state = SchedulerState()

local_req, local_strat, local_const = await env.sync_run_params(
    requests, strategy, constraints
)
start_time = await env.sync_run_start()
for req in local_req:
    state.processed_requests += 1
    await env.update_run_iteration(f"resp_{req}", req, RequestInfo(), state)
async for nonlocal_req in env.sync_run_end():
    state.processed_requests += 1
Source code in src/guidellm/scheduler/environments.py
class NonDistributedEnvironment(Environment[RequestT, ResponseT]):
    """
    Single-node scheduler execution environment with minimal coordination overhead.

    Implements the Environment interface with no-op synchronization for local testing,
    development, and single-machine benchmarking. All synchronization methods return
    immediately without distributed coordination logic.

    Example:
    ::
        from guidellm.scheduler import (
            MaxNumberConstraint,
            MaxRequestsConstraintArgs,
            NonDistributedEnvironment,
            RequestInfo,
            SchedulerState,
            SynchronousStrategy,
        )

        env = NonDistributedEnvironment()
        requests = [f"req_{ind}" for ind in range(5)]
        strategy = SynchronousStrategy()
        args = MaxRequestsConstraintArgs(max_num=5)
        constraints = {"max_requests": MaxNumberConstraint(args=args)}
        state = SchedulerState()

        local_req, local_strat, local_const = await env.sync_run_params(
            requests, strategy, constraints
        )
        start_time = await env.sync_run_start()
        for req in local_req:
            state.processed_requests += 1
            await env.update_run_iteration(f"resp_{req}", req, RequestInfo(), state)
        async for nonlocal_req in env.sync_run_end():
            state.processed_requests += 1
    """

    def __init__(self):
        """
        Initialize single-node environment with empty error storage.
        """
        self.run_errors: list[Exception] = []

    async def sync_run_params(
        self,
        requests: DatasetIterT[RequestT],
        strategy: SchedulingStrategy,
        constraints: dict[str, Constraint | ConstraintInitializer],
    ) -> tuple[
        DatasetIterT[RequestT],
        SchedulingStrategy,
        dict[str, Constraint | ConstraintInitializer],
    ]:
        """
        Return parameters unchanged for single-node execution.

        :param requests: Requests to process locally
        :param strategy: Scheduling strategy to apply during execution
        :param constraints: Runtime constraints to enforce during execution
        :return: Original (requests, strategy, constraints) tuple unchanged
        """
        return requests, strategy, constraints

    async def sync_run_start(self) -> float:
        """
        Return current time plus configured delay for single-node startup.

        :return: Unix timestamp when execution should begin
        """
        return time.time() + settings.scheduler_start_delay_non_distributed

    async def update_run_iteration(
        self,
        response: ResponseT | None,
        request: RequestT,
        request_info: RequestInfo,
        state: SchedulerState,
    ):
        """
        No-op for single-node execution with no distributed state synchronization.

        :param response: Response generated for the request, if successful
        :param request: The processed request
        :param request_info: Metadata about request processing including timings
        :param state: Current scheduler state with metrics and progress
        """

    async def sync_run_error(self, err: Exception | list[Exception]):
        """
        Store error for later propagation during run finalization.

        :param err: The exception(s) that occurred during execution
        """
        err = [err] if not isinstance(err, list) else err
        self.run_errors.extend(err)

    async def sync_run_end(
        self,
    ) -> AsyncIterator[
        tuple[
            ResponseT | None,
            RequestT,
            RequestInfo,
            SchedulerState,
        ]
    ]:
        """
        Finalize single-node execution and propagate any stored errors.

        :return: Empty iterator as there are no remote nodes
        :raises Exception: Any error stored during execution via sync_run_error
        """
        if self.run_errors:
            if len(self.run_errors) == 1:
                raise self.run_errors[0]
            else:
                raise RuntimeError(
                    f"Errors occurred during execution: {self.run_errors}"
                )

        if False:
            # Force compiler to recognize as generator
            yield None  # type: ignore[misc]

__init__()

Initialize single-node environment with empty error storage.

Source code in src/guidellm/scheduler/environments.py
def __init__(self):
    """
    Initialize single-node environment with empty error storage.
    """
    self.run_errors: list[Exception] = []

sync_run_end() async

Finalize single-node execution and propagate any stored errors.

Returns:

Type Description
AsyncIterator[tuple[ResponseT | None, RequestT, RequestInfo, SchedulerState]]

Empty iterator as there are no remote nodes

Raises:

Type Description
Exception

Any error stored during execution via sync_run_error

Source code in src/guidellm/scheduler/environments.py
async def sync_run_end(
    self,
) -> AsyncIterator[
    tuple[
        ResponseT | None,
        RequestT,
        RequestInfo,
        SchedulerState,
    ]
]:
    """
    Finalize single-node execution and propagate any stored errors.

    :return: Empty iterator as there are no remote nodes
    :raises Exception: Any error stored during execution via sync_run_error
    """
    if self.run_errors:
        if len(self.run_errors) == 1:
            raise self.run_errors[0]
        else:
            raise RuntimeError(
                f"Errors occurred during execution: {self.run_errors}"
            )

    if False:
        # Force compiler to recognize as generator
        yield None  # type: ignore[misc]

sync_run_error(err) async

Store error for later propagation during run finalization.

Parameters:

Name Type Description Default
err Exception | list[Exception]

The exception(s) that occurred during execution

required
Source code in src/guidellm/scheduler/environments.py
async def sync_run_error(self, err: Exception | list[Exception]):
    """
    Store error for later propagation during run finalization.

    :param err: The exception(s) that occurred during execution
    """
    err = [err] if not isinstance(err, list) else err
    self.run_errors.extend(err)

sync_run_params(requests, strategy, constraints) async

Return parameters unchanged for single-node execution.

Parameters:

Name Type Description Default
requests DatasetIterT[RequestT]

Requests to process locally

required
strategy SchedulingStrategy

Scheduling strategy to apply during execution

required
constraints dict[str, Constraint | ConstraintInitializer]

Runtime constraints to enforce during execution

required

Returns:

Type Description
tuple[DatasetIterT[RequestT], SchedulingStrategy, dict[str, Constraint | ConstraintInitializer]]

Original (requests, strategy, constraints) tuple unchanged

Source code in src/guidellm/scheduler/environments.py
async def sync_run_params(
    self,
    requests: DatasetIterT[RequestT],
    strategy: SchedulingStrategy,
    constraints: dict[str, Constraint | ConstraintInitializer],
) -> tuple[
    DatasetIterT[RequestT],
    SchedulingStrategy,
    dict[str, Constraint | ConstraintInitializer],
]:
    """
    Return parameters unchanged for single-node execution.

    :param requests: Requests to process locally
    :param strategy: Scheduling strategy to apply during execution
    :param constraints: Runtime constraints to enforce during execution
    :return: Original (requests, strategy, constraints) tuple unchanged
    """
    return requests, strategy, constraints

sync_run_start() async

Return current time plus configured delay for single-node startup.

Returns:

Type Description
float

Unix timestamp when execution should begin

Source code in src/guidellm/scheduler/environments.py
async def sync_run_start(self) -> float:
    """
    Return current time plus configured delay for single-node startup.

    :return: Unix timestamp when execution should begin
    """
    return time.time() + settings.scheduler_start_delay_non_distributed

update_run_iteration(response, request, request_info, state) async

No-op for single-node execution with no distributed state synchronization.

Parameters:

Name Type Description Default
response ResponseT | None

Response generated for the request, if successful

required
request RequestT

The processed request

required
request_info RequestInfo

Metadata about request processing including timings

required
state SchedulerState

Current scheduler state with metrics and progress

required
Source code in src/guidellm/scheduler/environments.py
async def update_run_iteration(
    self,
    response: ResponseT | None,
    request: RequestT,
    request_info: RequestInfo,
    state: SchedulerState,
):
    """
    No-op for single-node execution with no distributed state synchronization.

    :param response: Response generated for the request, if successful
    :param request: The processed request
    :param request_info: Metadata about request processing including timings
    :param state: Current scheduler state with metrics and progress
    """