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.
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).
OrchestratorRequestState dataclass ¶
Per-request bookkeeping inside the Orchestrator.
final_output_stage_ids class-attribute instance-attribute ¶
finished_final_output_stage_ids class-attribute instance-attribute ¶
finished_stage_ids class-attribute instance-attribute ¶
pd_prefill_multimodal_output class-attribute instance-attribute ¶
pending_final_output class-attribute instance-attribute ¶
pending_final_output: OutputMessage | None = None
pipeline_timings class-attribute instance-attribute ¶
request_artifact_dirs class-attribute instance-attribute ¶
running_counter_registered class-attribute instance-attribute ¶
running_counter_registered: bool = False
sampling_params_list class-attribute instance-attribute ¶
skip_legacy_stage_forward class-attribute instance-attribute ¶
skip_legacy_stage_forward: bool = False
stage_submit_ts class-attribute instance-attribute ¶
streaming class-attribute instance-attribute ¶
streaming: StreamingInputState = field(
default_factory=lambda: StreamingInputState()
)
StreamingInputState dataclass ¶
bridge_states class-attribute instance-attribute ¶
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 ¶
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.
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.