Shortcuts

Source code for torchx.runner.api

# Copyright (c) Meta Platforms, Inc. and affiliates.
# All rights reserved.
#
# This source code is licensed under the BSD-style license found in the
# LICENSE file in the root directory of this source tree.


import contextlib
import copy
import dataclasses
import json
import logging
import os
import time
import warnings
from datetime import datetime
from types import TracebackType
from typing import (
    TYPE_CHECKING,
    Any,
    Iterable,
    Iterator,
    Literal,
    Mapping,
    Type,
    TypeVar,
    overload,
)

from torchx import settings
from torchx.runner.events import log_event
from torchx.schedulers import SchedulerFactory, get_scheduler_factories
from torchx.schedulers.api import ListAppResponse, Scheduler, Stream
from torchx.specs import (
    AppDef,
    AppDryRunInfo,
    AppHandle,
    AppStatus,
    CfgVal,
    Role,
    UnknownAppException,
    UnknownSchedulerException,
    Workspace,
    macros,
    make_app_handle,
    materialize_appdef,
    parse_app_handle,
    runopts,
)
from torchx.specs.finder import get_component
from torchx.tracker.api import tracker_config_env_var_name
from torchx.util.session import get_session_id_or_create_new
from torchx.util.types import none_throws
from torchx.workspace import WorkspaceMixin

if TYPE_CHECKING:
    from typing_extensions import Self

from .config import get_config, get_configs

logger: logging.Logger = logging.getLogger(__name__)


NONE: str = "<NONE>"
S = TypeVar("S")
T = TypeVar("T")


@contextlib.contextmanager
def _overrides_detached(roles: list[Role]) -> Iterator[list[dict[str, Any]]]:
    """Strips ``overrides`` off *roles* for the duration, yielding the dicts.

    ``Role.overrides`` may hold non-deepcopyable values (APF attaches an
    in-flight fbpkg Future), so a ``deepcopy`` of anything owning a role runs
    with them detached -- as ``macros.Values.apply`` does. Attach the yielded
    dicts (the SAME objects) to the copy's roles: resolution writes back in
    place, so sharing lets whichever owner resolves first serve both.
    """
    detached = [role.overrides for role in roles]
    for role in roles:
        role.overrides = {}
    try:
        yield detached
    finally:
        for role, overrides in zip(roles, detached):
            role.overrides = overrides


def _built_only_the_image(authored: Role, built: Role) -> bool:
    """Whether building *authored*'s workspace changed nothing but ``image``.

    Builders that also write to the role (``env`` pointing at the unpacked
    workspace, say) produce a build the image alone cannot carry.

    Fields are compared with ``==``, so a field type without value equality
    can only turn this ``False`` -- costing a rebuild, never a bad pin.
    """
    return all(
        getattr(authored, f.name) == getattr(built, f.name)
        for f in dataclasses.fields(type(authored))
        if f.name != "image"
    )


def get_configured_trackers() -> dict[str, str | None]:
    tracker_names = list(get_configs(prefix="torchx", name="tracker").keys())
    if settings.ENV_TORCHX_TRACKERS in os.environ:
        tracker_names = os.environ[settings.ENV_TORCHX_TRACKERS].split(",")
        logger.info(
            "using %s=%s as tracker names", settings.ENV_TORCHX_TRACKERS, tracker_names
        )

    tracker_names_with_config = {}
    for tracker_name in tracker_names:
        config_value = get_config(prefix="tracker", name=tracker_name, key="config")

        config_env_name = tracker_config_env_var_name(tracker_name)
        if config_env_name in os.environ:
            config_value = os.environ[config_env_name]
            logger.info(
                "using %s=%s for `%s` tracker",
                config_env_name,
                config_value,
                tracker_name,
            )

        tracker_names_with_config[tracker_name] = config_value
    logger.info("tracker configurations: %s", tracker_names_with_config)
    return tracker_names_with_config


[docs] class Runner: """Submits, monitors, and manages :py:class:`~torchx.specs.AppDef` jobs. Use :py:func:`get_runner` to create an instance with all registered schedulers. .. doctest:: >>> from torchx.runner import get_runner >>> runner = get_runner() >>> runner.scheduler_backends() # doctest: +SKIP ['local_cwd', 'local_docker', 'slurm', 'kubernetes', ...] """ def __init__( self, name: str = "", # session names can be empty scheduler_factories: dict[str, SchedulerFactory] | None = None, component_defaults: dict[str, dict[str, str]] | None = None, scheduler_params: dict[str, object] | None = None, ) -> None: self._name: str = name self._scheduler_factories: dict[str, SchedulerFactory] = ( scheduler_factories if scheduler_factories is not None else get_scheduler_factories() ) self._scheduler_params: dict[str, Any] = { **(self._get_scheduler_params_from_env()), **(scheduler_params or {}), } self._scheduler_instances: dict[str, Scheduler] = {} self._apps: dict[AppHandle, AppDef] = {} # component_name -> map of component_fn_param_name -> user-specified default val encoded as str self._component_defaults: dict[str, dict[str, str]] = component_defaults or {} def _get_scheduler_params_from_env(self) -> dict[str, str]: scheduler_params = {} for key, value in os.environ.items(): key = key.lower() if key.startswith("torchx_"): scheduler_params[key.removeprefix("torchx_")] = value return scheduler_params # pyrefly: ignore [not-a-type] def __enter__(self) -> "Self": return self def __exit__( self, type: Type[BaseException] | None, value: BaseException | None, traceback: TracebackType | None, ) -> bool: # This method returns False so that if an error is raise within the # ``with`` statement, it is reraised properly # see: https://docs.python.org/3/reference/compound_stmts.html#with # see also: torchx/runner/test/api_test.py#test_context_manager_with_error # self.close() return False
[docs] def close(self) -> None: """Closes the runner and all scheduler instances. Safe to call multiple times.""" for scheduler in self._scheduler_instances.values(): scheduler.close()
[docs] def run_component( self, component: str, component_args: list[str] | dict[str, Any], scheduler: str, cfg: Mapping[str, CfgVal] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None, ) -> AppHandle: """Resolves and runs a named component. ``component`` resolution order (high → low): 1. User-registered ``torchx.components`` entry points 2. Builtins relative to ``torchx.components`` (e.g. ``"dist.torchrun"``) 3. File-based ``path/to/file.py:function_name`` """ with log_event("run_component") as ctx: dryrun_info = self.dryrun_component( component, component_args, scheduler, cfg=cfg, workspace=workspace, parent_run_id=parent_run_id, ) handle = self.schedule(dryrun_info) app = none_throws(dryrun_info.app) ctx._torchx_event.workspace = str(workspace) ctx._torchx_event.scheduler = none_throws(dryrun_info._scheduler) ctx._torchx_event.app_image = app.roles[0].image ctx._torchx_event.app_id = parse_app_handle(handle)[2] ctx._torchx_event.app_metadata = app.metadata return handle
[docs] def dryrun_component( self, component: str, component_args: list[str] | dict[str, Any], scheduler: str, cfg: Mapping[str, CfgVal] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None, ) -> AppDryRunInfo: """Like :py:meth:`run_component` but returns the request without submitting.""" component_def = get_component(component) args_from_cli = component_args if isinstance(component_args, list) else [] args_from_json = component_args if isinstance(component_args, dict) else {} app = materialize_appdef( component_def.fn, args_from_cli, self._component_defaults.get(component, None), args_from_json, ) return self.dryrun( app, scheduler, cfg=cfg, workspace=workspace, parent_run_id=parent_run_id, )
@overload def run( self, app: AppDef, scheduler: str, cfg: Mapping[str, CfgVal] | None = ..., workspace: Workspace | str | None = ..., parent_run_id: str | None = ..., *, dryrun: Literal[True], ) -> AppDryRunInfo: ... @overload def run( self, app: AppDef, scheduler: str, cfg: Mapping[str, CfgVal] | None = ..., workspace: Workspace | str | None = ..., parent_run_id: str | None = ..., *, dryrun: Literal[False] = ..., ) -> AppHandle: ...
[docs] def run( self, app: AppDef, scheduler: str, cfg: Mapping[str, CfgVal] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None, *, dryrun: bool = False, ) -> AppHandle | AppDryRunInfo: """Submits an :py:class:`~torchx.specs.AppDef` or returns its dry-run info. .. code-block:: python # Submit handle = runner.run(app, "mkube", cfg=cfg) # Dryrun — inspect without submitting info = runner.run(app, "mkube", cfg=cfg, dryrun=True) print(info) **Does not mutate** *app* -- see :py:meth:`dryrun`, which this delegates to. The workspace build and env injection land on an internal deep copy: with ``dryrun=True`` that copy comes back as ``info.app``; with ``dryrun=False`` it is submitted and dropped, and only the :py:data:`~torchx.specs.AppHandle` is returned. Args: dryrun: If ``True``, only validate and render the request without submitting. Returns :py:class:`~torchx.specs.AppDryRunInfo`. If ``False`` (default), submit and return the :py:data:`~torchx.specs.AppHandle`. Raises: ValueError: propagated from :py:meth:`dryrun`. """ with log_event(api="run") as ctx: dryrun_info = self.dryrun( app, scheduler, cfg=cfg, workspace=workspace, parent_run_id=parent_run_id, ) if dryrun: return dryrun_info handle = self.schedule(dryrun_info) event = ctx._torchx_event event.scheduler = scheduler event.runcfg = json.dumps(dict(cfg)) if cfg else None event.workspace = str(workspace) event.app_id = parse_app_handle(handle)[2] event.app_image = none_throws(dryrun_info.app).roles[0].image event.app_metadata = app.metadata return handle
[docs] def schedule(self, dryrun_info: AppDryRunInfo) -> AppHandle: """Submits a previously dry-run request, allowing request mutation. .. code-block:: python dryrun_info = runner.dryrun(app, scheduler="kubernetes", cfg) dryrun_info.request.foo = "bar" # mutate the raw request app_handle = runner.schedule(dryrun_info) .. warning:: Use sparingly. Overwriting many raw scheduler fields may cause your usage to diverge from TorchX's supported API. Only ``dryrun_info.request`` is submitted; edits to ``dryrun_info.app`` made after :py:meth:`dryrun` returned are **not** re-rendered into it. """ scheduler = none_throws(dryrun_info._scheduler) cfg = dryrun_info.cfg with log_event("schedule") as ctx: sched = self._scheduler(scheduler) app_id = sched.schedule(dryrun_info) app_handle = make_app_handle(scheduler, self._name, app_id) app = none_throws(dryrun_info.app) self._apps[app_handle] = app event = ctx._torchx_event event.scheduler = scheduler event.runcfg = json.dumps(dict(cfg)) if cfg else None event.app_id = app_id event.app_image = none_throws(dryrun_info.app).roles[0].image event.app_metadata = app.metadata return app_handle
def name(self) -> str: return self._name @overload def build_workspace( self, app_or_role: Role, scheduler: str, cfg: Mapping[str, CfgVal] | None = None, ) -> str | None: ... @overload def build_workspace( self, app_or_role: AppDef, scheduler: str, cfg: Mapping[str, CfgVal] | None = None, ) -> dict[tuple[str, Workspace], str]: ...
[docs] def build_workspace( self, app_or_role: AppDef | Role, scheduler: str, cfg: Mapping[str, CfgVal] | None = None, ) -> dict[tuple[str, Workspace], str] | str | None: """Builds the workspaces and returns the images they were built into. **Does not mutate** *app_or_role* -- the build runs against a deep copy. Use this to build once and submit the result many times, instead of letting each :py:meth:`run` rebuild: .. code-block:: python images = runner.build_workspace(app, "kubernetes", cfg) for app in per_region_apps: pin_workspace_images(app, images) runner.run(app, "kubernetes", cfg) Given a :py:class:`~torchx.specs.Role`, returns that role's built image, or ``None`` if it has no :py:attr:`~torchx.specs.Role.workspace`. Given an :py:class:`~torchx.specs.AppDef`, returns ``{(image, workspace): built}`` for the roles that have one -- keyed by the pair, not by role name, so the result can be pinned onto a *different* ``AppDef``. A build installs the workspace into ``role.image``, so the base image is half the key: roles sharing a workspace but not an image build to different results, and roles sharing both are built once. Both key halves come from the role as *authored*, which the build does not overwrite. The pair is the widest key correct for every builder, so one that caches more coarsely than the exact image (on the fbpkg name, say) may rebuild where it could have shared -- reuse is missed, never misapplied. Only :py:attr:`~torchx.specs.Role.workspace` is read. The deprecated ``workspace=`` argument to :py:meth:`run`/:py:meth:`dryrun` is applied inside ``dryrun``, too late to prebuild here; set the role attribute instead. Unlike :py:meth:`dryrun`, this neither validates *app_or_role* nor renders a scheduler request, so it does not fail on request-level problems that are unrelated to building. A role is omitted (with a warning) when its build wrote to the role beyond ``image`` -- ``env``, say. Pinning cannot replay those writes, so such a role keeps its workspace and rebuilds on submit. The whole ``(image, workspace)`` pair drops out, not just that role: a pair one role has to rebuild is not reusable for the roles that share it. Returns an empty mapping (or ``None``) when *scheduler* has no workspace support -- there is nothing to build. """ authored_roles = ( app_or_role.roles if isinstance(app_or_role, AppDef) else [app_or_role] ) with _overrides_detached(authored_roles) as detached_overrides: roles = copy.deepcopy(authored_roles) for role, overrides in zip(roles, detached_overrides): role.overrides = overrides with log_event("build_workspace", scheduler): sched = self._scheduler(scheduler) if not isinstance(sched, WorkspaceMixin): return None if isinstance(app_or_role, Role) else {} # one call, so roles sharing a workspace hit the shared build cache sched.build_workspaces(roles, sched.run_opts().resolve(cfg or {})) reusable: dict[tuple[str, Workspace], str] = {} rebuilt: set[tuple[str, Workspace]] = set() for authored, role in zip(authored_roles, roles, strict=True): workspace = authored.workspace if not workspace: continue if not _built_only_the_image(authored, role): logger.warning( "role `%s` cannot reuse its build: the `%s` workspace builder" " changed more than `image`, so the workspace has to be rebuilt" " on each submit", authored.name, scheduler, ) rebuilt.add((authored.image, workspace)) continue reusable[(authored.image, workspace)] = role.image for key in rebuilt: reusable.pop(key, None) if isinstance(app_or_role, Role): workspace = app_or_role.workspace return reusable.get((app_or_role.image, workspace)) if workspace else None return reusable
[docs] def dryrun( self, app: AppDef, scheduler: str, cfg: Mapping[str, CfgVal] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None, ) -> AppDryRunInfo: """Returns what *would* be submitted without actually submitting. The returned :py:class:`~torchx.specs.AppDryRunInfo` can be ``print()``-ed for inspection or passed to :py:meth:`schedule`. **Does not mutate** *app*. The patching this method performs -- building each role's :py:attr:`~torchx.specs.Role.workspace` and repointing ``role.image`` at the built artifact, injecting the ``TORCHX_*`` tracking env vars -- lands on a deep copy, which is returned as :py:attr:`~torchx.specs.AppDryRunInfo.app`. Read the submitted images off that copy, never off the *app* you passed in: .. code-block:: python app = AppDef(roles=[Role(image="foo:latest", workspace=..., ...)]) info = runner.dryrun(app, "kubernetes") app.roles[0].image # "foo:latest" -- unpatched, as authored info.app.roles[0].image # "foo:<built-workspace-hash>" -- submitted ``Role.overrides`` is the one part not copied: the same dict object is shared by the caller's role and the copy, so resolving an override through either is visible to both. Two dryruns of the same *app* need not render the same request -- each rebuilds the workspace, and the build may produce a new image. Re-running from the returned copy does not skip that build either, because :py:attr:`~torchx.specs.Role.workspace` is left set after one; it only narrows the difference down to the build, since every other patch is already applied. Clear the workspace to render the exact same request: .. code-block:: python info1 = runner.dryrun(app, "kubernetes", cfg) for role in info1.app.roles: role.workspace = None info2 = runner.dryrun(info1.app, "kubernetes", info1.cfg) assert info1.request == info2.request Raises: ValueError: *app* has no roles, or a role has no ``entrypoint`` or a non-positive ``num_replicas``. Schedulers raise from their own ``_validate`` hooks for backend-specific violations. """ # operate on a copy so that the env injection and workspace overwrite # below never leak into the caller's AppDef (one AppDef can be # dry-run multiple times); the copy is returned as the AppDryRunInfo's # `app`. with _overrides_detached(app.roles) as detached_overrides: app_copy = copy.deepcopy(app) for role_copy, overrides in zip(app_copy.roles, detached_overrides): role_copy.overrides = overrides app = app_copy # input validation if not app.roles: raise ValueError( f"No roles for app: {app.name}. Did you forget to add roles to AppDef?" ) if settings.ENV_TORCHX_PARENT_RUN_ID in os.environ: parent_run_id = os.environ[settings.ENV_TORCHX_PARENT_RUN_ID] logger.info( "using %s=%s env variable as tracker parent run id", settings.ENV_TORCHX_PARENT_RUN_ID, parent_run_id, ) configured_trackers = get_configured_trackers() for role in app.roles: if not role.entrypoint: raise ValueError( f"No entrypoint for role: {role.name}." f" Did you forget to call role.runs(entrypoint, args, env)?" ) if role.num_replicas <= 0: raise ValueError( f"Non-positive replicas for role: {role.name}." f" Did you forget to set role.num_replicas?" ) # Setup tracking # 1. Inject parent identifier # 2. Inject this run's job ID # 3. Get the list of backends to support from .torchconfig # - inject it as TORCHX_TRACKERS=names (it is expected that entrypoints are defined) # - for each backend check configuration file, if exists: # - inject it as TORCHX_TRACKER_<name>_CONFIGFILE=filename role.env[settings.ENV_TORCHX_JOB_ID] = make_app_handle( scheduler, self._name, macros.app_id ) role.env[settings.ENV_TORCHX_INTERNAL_SESSION_ID] = ( get_session_id_or_create_new() ) if parent_run_id: role.env[settings.ENV_TORCHX_PARENT_RUN_ID] = parent_run_id if configured_trackers: role.env[settings.ENV_TORCHX_TRACKERS] = ",".join( configured_trackers.keys() ) for name, config in configured_trackers.items(): if config: role.env[tracker_config_env_var_name(name)] = config cfg = cfg or dict() with log_event( "dryrun", scheduler, runcfg=json.dumps(dict(cfg)) if cfg else None, workspace=str(workspace), ) as ctx: sched = self._scheduler(scheduler) resolved_cfg = sched.run_opts().resolve(cfg) sched._pre_build_validate(app, scheduler, resolved_cfg) if isinstance(sched, WorkspaceMixin): if workspace: # NOTE: torchx originally took workspace as a runner arg and only applied the workspace to role[0] # later, torchx added support for the workspace attr in Role # for BC, give precedence to the workspace argument over the workspace attr for role[0] if app.roles[0].workspace: logger.info( "Overriding role[%d] (%s) workspace to `%s`" "To use the role's workspace attr pass: --workspace='' from CLI or workspace=None programmatically.", 0, role.name, str(app.roles[0].workspace), ) app.roles[0].workspace = ( Workspace.from_str(workspace) if isinstance(workspace, str) else workspace ) sched.build_workspaces(app.roles, resolved_cfg) sched._validate(app, scheduler, resolved_cfg) dryrun_info = sched.submit_dryrun(app, resolved_cfg) dryrun_info._scheduler = scheduler event = ctx._torchx_event event.scheduler = scheduler event.runcfg = json.dumps(dict(cfg)) if cfg else None event.app_id = app.name event.app_image = none_throws(dryrun_info.app).roles[0].image event.app_metadata = app.metadata return dryrun_info
[docs] def scheduler_run_opts(self, scheduler: str) -> runopts: """Returns the :py:class:`~torchx.specs.runopts` for the given scheduler.""" return self._scheduler(scheduler).run_opts()
[docs] def cfg_from_str(self, scheduler: str, *cfg_literal: str) -> Mapping[str, CfgVal]: """ Convenience function around the scheduler's ``runopts.cfg_from_str()`` method. Usage: .. doctest:: >>> from torchx.runner import get_runner >>> runner = get_runner() >>> runner.cfg_from_str("local_cwd", "log_dir=/tmp/foobar", "prepend_cwd=True") {'log_dir': '/tmp/foobar', 'prepend_cwd': True} """ opts = self._scheduler(scheduler).run_opts() cfg = {} for cfg_str in cfg_literal: cfg.update(opts.cfg_from_str(cfg_str)) return cfg
[docs] def scheduler_backends(self) -> list[str]: """Returns all registered scheduler backend names.""" return list(self._scheduler_factories.keys())
[docs] def status(self, app_handle: AppHandle) -> AppStatus | None: """Returns app status, or ``None`` if the app no longer exists.""" scheduler, scheduler_backend, app_id = self._scheduler_app_id( app_handle, check_session=False ) with log_event("status", scheduler_backend, app_id): desc = scheduler.describe(app_id) if not desc: # app does not exist on the scheduler # remove it from apps cache if it exists # effectively removes this app from the list() API self._apps.pop(app_handle, None) return None app_status = AppStatus( desc.state, desc.num_restarts, msg=desc.msg, structured_error_msg=desc.structured_error_msg, roles=desc.roles_statuses, ) if app_status: app_status.ui_url = desc.ui_url return app_status
[docs] def wait( self, app_handle: AppHandle, wait_interval: float = 10 ) -> AppStatus | None: """Blocks until the app reaches a terminal state. Args: wait_interval: seconds between status polls """ scheduler, scheduler_backend, app_id = self._scheduler_app_id( app_handle, check_session=False ) with log_event("wait", scheduler_backend, app_id): while True: app_status = self.status(app_handle) if not app_status: return None if app_status.is_terminal(): return app_status else: time.sleep(wait_interval)
[docs] def cancel(self, app_handle: AppHandle) -> None: """Requests cancellation. The app transitions to ``CANCELLED`` asynchronously.""" scheduler, scheduler_backend, app_id = self._scheduler_app_id(app_handle) with log_event("cancel", scheduler_backend, app_id): status = self.status(app_handle) if status is not None and not status.is_terminal(): scheduler.cancel(app_id)
[docs] def delete(self, app_handle: AppHandle) -> None: """Deletes the app from the scheduler.""" scheduler, scheduler_backend, app_id = self._scheduler_app_id(app_handle) with log_event("delete", scheduler_backend, app_id): status = self.status(app_handle) if status is not None: scheduler.delete(app_id)
[docs] def stop(self, app_handle: AppHandle) -> None: """.. deprecated:: Use :py:meth:`cancel` instead.""" warnings.warn( "This method will be deprecated in the future, please use `cancel` instead.", PendingDeprecationWarning, ) self.cancel(app_handle)
[docs] def describe(self, app_handle: AppHandle) -> AppDef | None: """Reconstructs the :py:class:`~torchx.specs.AppDef` from the scheduler. Completeness is scheduler-dependent. Returns ``None`` if the app no longer exists. """ scheduler, scheduler_backend, app_id = self._scheduler_app_id( app_handle, check_session=False ) with log_event("describe", scheduler_backend, app_id): # if the app is in the apps list, then short circuit everything and return it app = self._apps.get(app_handle, None) if not app: desc = scheduler.describe(app_id) if desc: app = AppDef(name=app_id, roles=desc.roles, metadata=desc.metadata) return app
[docs] def describe_native(self, app_handle: AppHandle) -> AppDryRunInfo | None: """Reads the scheduler-native request back for an already-submitted app. The read-side twin of :py:meth:`dryrun`: same :py:class:`~torchx.specs.AppDryRunInfo` wrapper, same scheduler-specific ``request`` type. Unlike :py:meth:`describe`, the returned ``request`` preserves scheduler-specific fields that have no :py:class:`~torchx.specs.AppDef` equivalent. Returns ``None`` if the scheduler does not support native read-back or the app no longer exists. """ scheduler, scheduler_backend, app_id = self._scheduler_app_id( app_handle, check_session=False ) with log_event("describe_native", scheduler_backend, app_id): info = scheduler.describe_native(app_id) if info is not None: info._scheduler = scheduler_backend return info
[docs] def log_lines( self, app_handle: AppHandle, role_name: str, k: int = 0, regex: str | None = None, since: datetime | None = None, until: datetime | None = None, should_tail: bool = False, streams: Stream | None = None, ) -> Iterable[str]: """Returns an iterator over log lines for the k-th replica of a role. .. important:: ``k`` is the **node** (host) id, NOT the worker rank. .. warning:: Completeness is scheduler-dependent. Lines may be partial or missing if logs have been purged. Do not use this for programmatic output parsing. Lines include trailing whitespace (``\\n``). Use ``print(line, end="")`` to avoid double newlines. Args: k: replica (node) index regex: optional filter pattern since: start cursor (scheduler-dependent) until: end cursor (scheduler-dependent) """ scheduler, scheduler_backend, app_id = self._scheduler_app_id( app_handle, check_session=False ) with log_event("log_lines", scheduler_backend, app_id): if not self.status(app_handle): raise UnknownAppException(app_handle) log_iter = scheduler.log_iter( app_id, role_name, k, regex, since, until, should_tail, streams=streams, ) return log_iter
[docs] def list( self, scheduler: str, cfg: Mapping[str, CfgVal] | None = None, ) -> list[ListAppResponse]: """Lists jobs on the scheduler. Args: cfg: scheduler config, used by some schedulers for backend routing. """ with log_event("list", scheduler): sched = self._scheduler(scheduler) apps = sched.list(cfg) for app in apps: app.app_handle = make_app_handle(scheduler, self._name, app.app_id) return apps
def _scheduler(self, scheduler: str) -> Scheduler: sched = self._scheduler_instances.get(scheduler) if not sched: factory = self._scheduler_factories.get(scheduler) if factory: sched = factory(self._name, **self._scheduler_params) self._scheduler_instances[scheduler] = sched if not sched: raise UnknownSchedulerException(scheduler) return sched def _scheduler_app_id( self, app_handle: AppHandle, check_session: bool = True, ) -> tuple[Scheduler, str, str]: scheduler_backend, _, app_id = parse_app_handle(app_handle) scheduler = self._scheduler(scheduler_backend) return scheduler, scheduler_backend, app_id def __repr__(self) -> str: return f"Runner(name={self._name}, schedulers={self._scheduler_factories}, apps={self._apps})"
[docs] def get_runner( name: str | None = None, component_defaults: dict[str, dict[str, str]] | None = None, **scheduler_params: Any, ) -> Runner: """Creates a :py:class:`Runner` with all registered schedulers. .. code-block:: python with get_runner() as runner: app_handle = runner.run(app, scheduler="kubernetes", cfg=cfg) print(runner.status(app_handle)) Args: scheduler_params: extra kwargs passed to all scheduler constructors. """ if name: warnings.warn( f"Custom session names are deprecated (detected explicitly set session name={name}). \ To prevent this warning from showing again call `get_runner()` without the `name` param. \ As an alternative, you can prefix the app name with the session name.", FutureWarning, ) if not name: name = "torchx" scheduler_factories = get_scheduler_factories() return Runner( name, scheduler_factories, component_defaults, scheduler_params=scheduler_params )

Docs

Access comprehensive developer documentation for PyTorch

View Docs

Tutorials

Get in-depth tutorials for beginners and advanced developers

View Tutorials

Resources

Find development resources and get your questions answered

View Resources