Skip to content

vllm_omni.diffusion.cancellation

Cooperative cancellation for local, full-forward diffusion workers.

The engine owns each signal until the worker has returned. Unlike an executor RPC, writing the signal does not wait behind the forward being cancelled. Opted-in pipelines call check_request_cancellation at safe boundaries.

RequestCancellationRegistry

Engine-owned signals; scheduler mutation remains on the engine thread.

cancel

cancel(request_ids: Iterable[str]) -> None

cancel_all

cancel_all() -> None

close

close() -> None

Release remaining signals after the executor has shut down.

create

create(request_id: str) -> str

finish

finish(request_id: str) -> None

Release only after execution returns, including aborted execution.

check_request_cancellation

check_request_cancellation(
    *, synchronize: bool = False
) -> None

Stop a cancelled execution wave without stranding a peer's collectives.

synchronize=True drains queued device work only when cancellation has been requested locally, before re-reading the flags. Every abort drains queued device work before unwinding. Successful steps retain asynchronous execution when no local cancellation is observed. Independent requests coupled by an AllGather offload wave must all be cancelled before the wave can exit; cancelling one must not abort its live peers.

request_cancellation_scope

request_cancellation_scope(
    signal_names: Sequence[str | None],
    *,
    enabled: bool = True,
) -> Iterator[None]

Attach worker readers without taking ownership of the engine's names.