Skip to content

vllm_omni.engine.orchestrator

Orchestrator for vLLM-Omni multi-stage runtime.

Runs inside a background thread with its own asyncio event loop. Owns logical request progression across stage pools and handles stage-to-stage transfer logic.

Distributed membership (replica attach/detach, hub monitoring) is handled by :class:MembershipController, which is injected optionally.

logger module-attribute

logger = init_logger(__name__)

Orchestrator

Bases: OrchestratorBase

Turn-based orchestrator: admits add_request / streaming / companion / interaction messages.

OrchestratorBase

Stage-management loop shared by the turn-based and duplex orchestrators.

Runs inside a background thread's asyncio event loop. Subclasses fill the template seams (_dispatch_message, _background_tasks, _shutdown_extensions, _on_stage_submitted, _intercept_stage_output, _handle_forward_failure, _cleanup_request_ids).

async_chunk instance-attribute

async_chunk = bool(async_chunk)

log_stats instance-attribute

log_stats = log_stats

num_stages instance-attribute

num_stages = len(stage_pools)

output_async_queue instance-attribute

output_async_queue = output_async_queue

request_async_queue instance-attribute

request_async_queue = request_async_queue

request_states instance-attribute

request_states: dict[str, OrchestratorRequestState] = {}

rpc_async_queue instance-attribute

rpc_async_queue = rpc_async_queue

stage_pools instance-attribute

stage_pools: list[StagePool] = stage_pools

run async

run() -> None

Main entry point for the Orchestrator event loop.

OrchestratorRequestState dataclass

Per-request bookkeeping inside the Orchestrator.

final_output_stage_ids class-attribute instance-attribute

final_output_stage_ids: set[int] = field(
    default_factory=set
)

final_stage_id class-attribute instance-attribute

final_stage_id: int = -1

finished_final_output_stage_ids class-attribute instance-attribute

finished_final_output_stage_ids: set[int] = field(
    default_factory=set
)

finished_stage_ids class-attribute instance-attribute

finished_stage_ids: set[int] = field(default_factory=set)

mm_features class-attribute instance-attribute

mm_features: list | None = None

mm_processor_kwargs class-attribute instance-attribute

mm_processor_kwargs: dict | None = None

native_kv_transfer_id class-attribute instance-attribute

native_kv_transfer_id: str | None = None

pd_prefill_multimodal_output class-attribute instance-attribute

pd_prefill_multimodal_output: dict[str, Any] | None = None

pending_final_output class-attribute instance-attribute

pending_final_output: OutputMessage | None = None

pipeline_timings class-attribute instance-attribute

pipeline_timings: dict[str, float] = field(
    default_factory=dict
)

prompt class-attribute instance-attribute

prompt: Any = None

request_artifact_dirs class-attribute instance-attribute

request_artifact_dirs: set[str] = field(default_factory=set)

request_id instance-attribute

request_id: str

request_timestamp class-attribute instance-attribute

request_timestamp: float = 0.0

running_counter_registered class-attribute instance-attribute

running_counter_registered: bool = False

sampling_params_list class-attribute instance-attribute

sampling_params_list: list[Any] = field(
    default_factory=list
)

session_owned class-attribute instance-attribute

session_owned: bool = False

skip_legacy_stage_forward class-attribute instance-attribute

skip_legacy_stage_forward: bool = False

stage_submit_ts class-attribute instance-attribute

stage_submit_ts: dict[int, float] = field(
    default_factory=dict
)

streaming class-attribute instance-attribute

streaming: StreamingInputState = field(
    default_factory=lambda: StreamingInputState()
)

StreamingInputState dataclass

bridge_states class-attribute instance-attribute

bridge_states: dict[str, Any] = field(default_factory=dict)

enabled class-attribute instance-attribute

enabled: bool = False

new_prompt_len_snapshot class-attribute instance-attribute

new_prompt_len_snapshot: int | None = None

segments class-attribute instance-attribute

segments: dict[int, StreamingSegmentState] = field(
    default_factory=dict
)

source_token_decoder class-attribute instance-attribute

source_token_decoder: Callable[..., str] | None = None

segment

segment(stage_id: int | None) -> StreamingSegmentState

Return stage_id's segment state, or a fresh empty one if unreported.

Fresh rather than shared: output_metadata is handed out by reference.

StreamingSegmentState dataclass

Streaming segment-boundary state for one stage of a request.

Every stage of a multi-stage pipeline reaches its own segment boundaries independently, so this is tracked per stage rather than per request.

finished class-attribute instance-attribute

finished: bool = False

output_metadata class-attribute instance-attribute

output_metadata: dict[str, Any] = field(
    default_factory=dict
)

token_ids class-attribute instance-attribute

token_ids: list[int] = field(default_factory=list)

build_engine_core_request_from_tokens

build_engine_core_request_from_tokens(
    request_id: str,
    prompt: dict[str, Any],
    params: SamplingParams | PoolingParams,
    arrival_time: float | None = None,
    model_config: ModelConfig | None = None,
    resumable: bool = False,
    mm_features: list | None = None,
) -> OmniEngineCoreRequest

Build an OmniEngineCoreRequest directly from an OmniTokensPrompt.

cleanup_request_artifact_dirs

cleanup_request_artifact_dirs(
    artifact_dirs: set[str] | list[str],
) -> None