Skip to content

external_launcher_executor

In-process executor for externally launched SPMD inference.

Unlike :class:MultiprocExecutor, which starts local subprocesses, this executor turns each process started by torchrun or srun into one inline worker. The launcher owns RANK, WORLD_SIZE, LOCAL_RANK, MASTER_ADDR, and MASTER_PORT; FastVideo initializes through the env:// rendezvous.

Every rank must call generation collectively with the same requests in the same order. World rank 0 owns saved media and returned frames; other ranks still execute the full pipeline so sequence-parallel collectives stay uniform.

Classes

fastvideo.worker.external_launcher_executor.ExternalLauncherEnv dataclass

ExternalLauncherEnv(rank: int, local_rank: int, world_size: int, local_world_size: int | None, master_addr: str, master_port: int)

Distributed identity assigned to this process by its launcher.

fastvideo.worker.external_launcher_executor.ExternalLauncherExecutor

ExternalLauncherExecutor(fastvideo_args: FastVideoArgs, *, log_queue=None)

Bases: Executor

Run this externally launched process's single worker inline.

Source code in fastvideo/worker/executor.py
def __init__(
    self,
    fastvideo_args: FastVideoArgs,
    *,
    log_queue=None,
):
    self.fastvideo_args = fastvideo_args
    self._log_queue = log_queue

    self._init_executor()

Methods:

fastvideo.worker.external_launcher_executor.ExternalLauncherExecutor.broadcast_from_output_rank
broadcast_from_output_rank(value)

Broadcast a small control-plane value from global rank zero.

Source code in fastvideo/worker/external_launcher_executor.py
def broadcast_from_output_rank(self, value):
    """Broadcast a small control-plane value from global rank zero."""
    return get_world_group().broadcast_object(value, src=0)
fastvideo.worker.external_launcher_executor.ExternalLauncherExecutor.collective_rpc
collective_rpc(method: str | Callable, timeout: float | None = None, args: tuple = (), kwargs: dict | None = None) -> list[Any]

Execute a control call locally and return rank-ordered world results.

Every launcher process must enter this method with the same call. A local exception is gathered before any rank raises, so ordinary control-plane failures are reported consistently instead of leaving a successful peer behind. Per-call timeouts are unsupported; the process-group timeout and launcher failure policy bound hard failures.

Source code in fastvideo/worker/external_launcher_executor.py
def collective_rpc(
        self,
        method: str | Callable,
        timeout: float | None = None,
        args: tuple = (),
        kwargs: dict | None = None,
) -> list[Any]:
    """Execute a control call locally and return rank-ordered world results.

    Every launcher process must enter this method with the same call. A
    local exception is gathered before any rank raises, so ordinary
    control-plane failures are reported consistently instead of leaving a
    successful peer behind. Per-call timeouts are unsupported; the
    process-group timeout and launcher failure policy bound hard failures.
    """

    if timeout is not None:
        raise NotImplementedError(
            "ExternalLauncherExecutor.collective_rpc does not support per-call timeouts; "
            "set generator.engine.parallelism.dist_timeout in seconds and configure the launcher's "
            "kill-on-failure policy instead.")
    self._bind_worker_device()
    try:
        response = self.worker.execute_method(method, *args, **(kwargs or {}))
        local_status: dict[str, Any] = {
            "ok": True,
            "rank": self.rank,
            "response": response,
        }
    except Exception as error:
        local_status = {
            "ok": False,
            "rank": self.rank,
            "error_type": type(error).__name__,
            "error": str(error),
        }

    statuses: list[dict[str, Any] | None]
    if self.world_size == 1:
        statuses = [local_status]
    else:
        statuses = [None] * self.world_size
        torch.distributed.all_gather_object(statuses, local_status, group=get_world_group().cpu_group)

    failures = [status for status in statuses if status is not None and not status["ok"]]
    if failures:
        details = "; ".join(f"rank {failure['rank']} {failure['error_type']}: {failure['error']}"
                            for failure in failures)
        method_name = method if isinstance(method, str) else getattr(method, "__name__", repr(method))
        raise RuntimeError(f"Collective RPC {method_name!r} failed: {details}")
    return [status["response"] for status in statuses if status is not None]

Functions:

fastvideo.worker.external_launcher_executor.read_external_launcher_environ

read_external_launcher_environ() -> dict[str, str]

Return the launcher variables that :func:resolve_external_launcher_env reads.

Each name is read explicitly rather than handing over the whole process environment (docs/contributing/env_vars.md): torchrun's identity variables and Slurm's per-task SLURM_* identity for native srun launches.

Source code in fastvideo/worker/external_launcher_executor.py
def read_external_launcher_environ() -> dict[str, str]:
    """Return the launcher variables that :func:`resolve_external_launcher_env` reads.

    Each name is read explicitly rather than handing over the whole process
    environment (docs/contributing/env_vars.md): torchrun's identity variables
    and Slurm's per-task ``SLURM_*`` identity for native ``srun`` launches.
    """
    values = {
        "RANK": os.environ.get("RANK"),
        "WORLD_SIZE": os.environ.get("WORLD_SIZE"),
        "LOCAL_RANK": os.environ.get("LOCAL_RANK"),
        "LOCAL_WORLD_SIZE": os.environ.get("LOCAL_WORLD_SIZE"),
        "MASTER_ADDR": os.environ.get("MASTER_ADDR"),
        "MASTER_PORT": os.environ.get("MASTER_PORT"),
        "SLURM_PROCID": os.environ.get("SLURM_PROCID"),
        "SLURM_NTASKS": os.environ.get("SLURM_NTASKS"),
        "SLURM_LOCALID": os.environ.get("SLURM_LOCALID"),
        "SLURM_NTASKS_PER_NODE": os.environ.get("SLURM_NTASKS_PER_NODE"),
        "SLURM_STEP_TASKS_PER_NODE": os.environ.get("SLURM_STEP_TASKS_PER_NODE"),
        "SLURM_TASKS_PER_NODE": os.environ.get("SLURM_TASKS_PER_NODE"),
    }
    return {name: value for name, value in values.items() if value is not None}

fastvideo.worker.external_launcher_executor.resolve_external_launcher_env

resolve_external_launcher_env(environ: Mapping[str, str]) -> ExternalLauncherEnv

Parse and validate a torchrun/srun distributed environment.

Source code in fastvideo/worker/external_launcher_executor.py
def resolve_external_launcher_env(environ: Mapping[str, str]) -> ExternalLauncherEnv:
    """Parse and validate a torchrun/srun distributed environment."""

    rank_value, rank_source = _first_env_value(environ, "RANK", "SLURM_PROCID")
    world_size_value, world_size_source = _first_env_value(environ, "WORLD_SIZE", "SLURM_NTASKS")
    master_addr, _ = _first_env_value(environ, "MASTER_ADDR")
    master_port_value, _ = _first_env_value(environ, "MASTER_PORT")
    missing = []
    if rank_value is None:
        missing.append("RANK (or SLURM_PROCID)")
    if world_size_value is None:
        missing.append("WORLD_SIZE (or SLURM_NTASKS)")
    if master_addr is None:
        missing.append("MASTER_ADDR")
    if master_port_value is None:
        missing.append("MASTER_PORT")
    if missing:
        raise RuntimeError("External-launcher inference requires the launcher to provide "
                           f"{missing}. Launch with torchrun/srun and an env:// rendezvous.")

    assert rank_value is not None and rank_source is not None
    assert world_size_value is not None and world_size_source is not None
    assert master_addr is not None and master_port_value is not None
    rank = _parse_int(rank_source, rank_value)
    world_size = _parse_int(world_size_source, world_size_value)
    master_port = _parse_int("MASTER_PORT", master_port_value)

    local_rank_value, local_rank_source = _first_env_value(environ, "LOCAL_RANK", "SLURM_LOCALID")
    if local_rank_value is None:
        if world_size == 1:
            local_rank = 0
        else:
            raise RuntimeError("External-launcher inference requires LOCAL_RANK (torchrun) "
                               "or SLURM_LOCALID (srun) so each process can select its GPU.")
    else:
        assert local_rank_source is not None
        local_rank = _parse_int(local_rank_source, local_rank_value)

    local_world_size_value, local_world_size_source = _first_env_value(environ, "LOCAL_WORLD_SIZE")
    local_world_size: int | None
    if local_world_size_value is not None:
        assert local_world_size_source is not None
        local_world_size = _parse_int(local_world_size_source, local_world_size_value)
    else:
        local_world_size = _slurm_local_world_size(environ)

    if world_size < 1:
        raise ValueError(f"WORLD_SIZE must be >= 1, got {world_size}.")
    if not 0 <= rank < world_size:
        raise ValueError(f"RANK must be in [0, WORLD_SIZE); got RANK={rank}, WORLD_SIZE={world_size}.")
    if local_rank < 0:
        raise ValueError(f"LOCAL_RANK must be >= 0, got {local_rank}.")
    if local_world_size is not None:
        if local_world_size < 1:
            raise ValueError(f"LOCAL_WORLD_SIZE must be >= 1, got {local_world_size}.")
        if local_rank >= local_world_size:
            raise ValueError("LOCAL_RANK must be in [0, LOCAL_WORLD_SIZE); "
                             f"got LOCAL_RANK={local_rank}, LOCAL_WORLD_SIZE={local_world_size}.")
    elif local_rank >= world_size:
        raise ValueError(
            f"LOCAL_RANK must be in [0, WORLD_SIZE); got LOCAL_RANK={local_rank}, WORLD_SIZE={world_size}.")
    if not 1 <= master_port <= 65535:
        raise ValueError(f"MASTER_PORT must be in [1, 65535], got {master_port}.")
    return ExternalLauncherEnv(
        rank=rank,
        local_rank=local_rank,
        world_size=world_size,
        local_world_size=local_world_size,
        master_addr=master_addr,
        master_port=master_port,
    )