Skip to content

uniproc_executor

In-process executor: one worker in the current process, no spawn.

Classes

fastvideo.worker.uniproc_executor.UniprocExecutor

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

Bases: Executor

Run a single Worker in the current process.

Used when num_gpus == 1 (including the default mp backend) or when distributed_executor_backend == "uni". Weights load once; no child process is spawned.

Because the worker shares the caller's process there is no workers list and no worker-side SIGINT handler: Ctrl-C reaches the main thread only, so call :meth:interrupt to cancel the in-flight forward pass. Cancellation is scoped to a generation (:meth:begin_generation): a request that lands while no generation is open is dropped, and a cancelled run raises RuntimeError instead of returning undenoised output as success.

StreamingVideoGenerator(use_queue_mode=True) prefetching needs a MultiprocExecutor; with this executor it warns and falls back to per-step in-process calls.

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.uniproc_executor.UniprocExecutor.begin_generation
begin_generation() -> None

Mark the start of a generation.

Drops any stale cancellation left over from a previous run so it can never silently skip this one, and opens the window in which :meth:interrupt is honored (closed again when the run ends).

Source code in fastvideo/worker/uniproc_executor.py
def begin_generation(self) -> None:
    """Mark the start of a generation.

    Drops any stale cancellation left over from a previous run so it can
    never silently skip this one, and opens the window in which
    :meth:`interrupt` is honored (closed again when the run ends).
    """
    self._generation_open = True
    self._clear_interrupt()
fastvideo.worker.uniproc_executor.UniprocExecutor.collective_rpc
collective_rpc(method: str | Callable, timeout: float | None = None, args: tuple = (), kwargs: dict | None = None) -> list[Any]

Execute the method on the in-process worker.

timeout is accepted for Executor compatibility but ignored: the call runs on the caller's thread, so a hung call blocks indefinitely.

Source code in fastvideo/worker/uniproc_executor.py
def collective_rpc(self,
                   method: str | Callable,
                   timeout: float | None = None,
                   args: tuple = (),
                   kwargs: dict | None = None) -> list[Any]:
    """Execute the method on the in-process worker.

    ``timeout`` is accepted for Executor compatibility but ignored: the
    call runs on the caller's thread, so a hung call blocks indefinitely.
    """
    del timeout
    kwargs = kwargs or {}
    return [self.driver_worker.execute_method(method, *args, **kwargs)]
fastvideo.worker.uniproc_executor.UniprocExecutor.interrupt
interrupt() -> None

Request cancellation of the in-flight forward pass (best effort).

There is no worker process to signal, so flag the in-process pipeline stages whose denoising loop checks self.interrupt instead. The request is honored only while a generation is open (from :meth:begin_generation, or the first call, until the run ends): a cancel that lands while no generation is open is dropped, and a cancelled run raises RuntimeError instead of returning undenoised output as success.

Source code in fastvideo/worker/uniproc_executor.py
def interrupt(self) -> None:
    """Request cancellation of the in-flight forward pass (best effort).

    There is no worker process to signal, so flag the in-process pipeline
    stages whose denoising loop checks ``self.interrupt`` instead. The
    request is honored only while a generation is open (from
    :meth:`begin_generation`, or the first call, until the run ends): a
    cancel that lands while no generation is open is dropped, and a
    cancelled run raises ``RuntimeError`` instead of returning undenoised
    output as success.
    """
    if not getattr(self, "_generation_open", False):
        return
    for stage in self._pipeline_stages():
        stage.interrupt = True
fastvideo.worker.uniproc_executor.UniprocExecutor.set_log_queue
set_log_queue(log_queue: Queue | None) -> None

Forward in-process logs to the given queue.

Unlike MultiprocExecutor, the handler is attached to this (driver) process's fastvideo logger, not to a worker, so driver logs reach the queue as well.

Source code in fastvideo/worker/uniproc_executor.py
def set_log_queue(self, log_queue: Queue | None) -> None:
    """Forward in-process logs to the given queue.

    Unlike MultiprocExecutor, the handler is attached to this (driver)
    process's ``fastvideo`` logger, not to a worker, so driver logs reach
    the queue as well.
    """
    self._clear_log_queue_handler()
    self._log_queue = log_queue
    if log_queue is None:
        return
    self._log_queue_handler = _make_queue_log_handler(log_queue)
    logging.getLogger("fastvideo").addHandler(self._log_queue_handler)

Functions: