Skip to content

vllm_omni.engine.duplex.plugin

Model plugin contract for full-duplex models.

One DuplexModelPlugin subclass per model binds what used to be two dotted paths (the engine DuplexRuntimeExtension and the serving ServingRuntimeAdapter). Everything runs engine-side now, so the plugin is loaded once by DuplexOmniEngine and handed to DuplexOrchestrator / DuplexSessionManager.

EncodeAudio module-attribute

EncodeAudio = Callable[
    [object, int, str, float | None], str | None
]

DuplexDataPlane

Bases: ABC

Projects raw stage outputs of one model into internal duplex events.

begin_request abstractmethod

begin_request(request_id: str) -> None

close_session abstractmethod

close_session(
    session_id: str, *, active_request_id: str | None = None
) -> None

close_stream abstractmethod

close_stream(request_id: str) -> None

is_terminal abstractmethod

is_terminal(request_id: str | None) -> bool

mark_terminal abstractmethod

mark_terminal(request_id: str) -> None

project abstractmethod

project(
    result: object, *, context: object | None = None
) -> Iterable[dict[str, object]]

DuplexModelPlugin

Bases: ABC

Everything vLLM-Omni needs to know about one full-duplex model.

Engine policy (sampling params, append planning, output decisions) and session policy (capabilities, runtime configuration, per-session state, data-plane projection) live on the same object so a mismatch between the two halves is impossible by construction.

data_plane instance-attribute

data_plane: DuplexDataPlane

plugin_id class-attribute instance-attribute

plugin_id: str = ''

private_runtime_config_keys class-attribute instance-attribute

private_runtime_config_keys: frozenset[str] = frozenset()

projects_intermediate_outputs class-attribute instance-attribute

projects_intermediate_outputs: bool = False

silence_continuation_samples class-attribute instance-attribute

silence_continuation_samples: int = 16000

capabilities abstractmethod

capabilities(*, max_sessions: int) -> DuplexCapabilities

commit_model_context

commit_model_context(
    *, session_id: str | None, assistant_text: str
) -> None

Persist model-context history at a turn boundary. Default is a no-op.

This is not playback-ACK history. A model that keeps its own prompt transcript implements this; the session runner only decides when.

configure_sampling_params abstractmethod

configure_sampling_params(
    *,
    runtime_config: dict[str, object],
    defaults: tuple[object, ...],
) -> tuple[object, ...]

create_session_state abstractmethod

create_session_state() -> DuplexModelSessionState

data_plane_context abstractmethod

data_plane_context(
    *,
    epoch: int,
    turn_id: int,
    active_response_turn_id: int | None,
    active_response_id: str | None,
    auto_responds: bool,
    response_format: str,
    speed: float | None,
    modalities: tuple[str, ...],
) -> object

decide_output abstractmethod

decide_output(
    *,
    stage_id: int,
    final_stage_id: int,
    segment_finished: bool,
    segment_token_ids: tuple[int, ...],
    segment_output_metadata: dict[str, object],
    output: object,
) -> DuplexOutputDecision | None

draining_stage_ids

draining_stage_ids(*, stage_count: int) -> frozenset[int]

Output stages that may keep running after the next user turn starts.

Empty means a concurrent turn does not overlap a previous output stage. Shared lifecycle reads this instead of assuming a stage layout.

partial_stage_followup

partial_stage_followup(
    plan: PartialStageForward, req_state: object
) -> PartialStageForward | None

Optional second submit after plan has already been forwarded.

Default models have nothing to add. AURA uses this to queue the end sentinel only after the last sentence text is already resumable.

plan_append abstractmethod

plan_append(
    *,
    request_id: str,
    fence: DuplexFence,
    session_config: dict[str, object],
    runtime_config: dict[str, object],
    seq: int,
    turn_seq: int,
    payload: object,
    final: bool,
    sampling_params: object,
) -> DuplexAppendPlan

plan_partial_stage_output

plan_partial_stage_output(
    orchestrator: object,
    stage_id: int,
    replica_id: int,
    output: object,
    req_state: object,
) -> PartialStageForward | None

Return a Talker update the orchestrator should submit, or None.

Default models do not split Stage1 text. The orchestrator owns the actual _forward_to_next_stage call.

prepare_append_plan async

prepare_append_plan(**kwargs) -> DuplexAppendPlan

Prepare a plan; plugins may offload expensive work on owned snapshots.

prepare_prompt_config

prepare_prompt_config(
    config: dict[str, object],
    *,
    state: DuplexModelSessionState,
    payload: dict[str, object],
) -> dict[str, object]

Add model-owned context before planning an append on the session loop.

prepare_runtime_config abstractmethod async

prepare_runtime_config(
    config: DuplexSessionConfig,
    *,
    model_config: ModelConfig | None,
) -> dict[str, object]

project_intermediate_output

project_intermediate_output(
    *, stage_id: int, output: object, context: object
) -> bool

Return True to project this intermediate stage to the client.

Unlike decide_output, projecting does not short-circuit the pipeline: the stage output is still forwarded to the next stage. Default is off. Orthogonal to projects_intermediate_outputs (Qwen3 Stage0); this hook is per-stage.

release_concurrent_turn_requests

release_concurrent_turn_requests(
    *,
    stage_id: int,
    segment_finished: bool,
    output: object,
    context: object,
) -> bool

Return True when the next user commit may start while prior TTS drains.

The plugin chooses when that is safe. The runner must not hard-code a stage id. Default off. Unlike barge-in, this path must not cancel the old TTS.

runtime_config_for_function_output

runtime_config_for_function_output(
    config: DuplexSessionConfig,
    current: Mapping[str, object],
    item: Mapping[str, object],
) -> dict[str, object] | None

runtime_config_for_update abstractmethod

runtime_config_for_update(
    config: DuplexSessionConfig,
    current: Mapping[str, object],
) -> dict[str, object]

user_transcript

user_transcript(
    *,
    stage_id: int,
    output: object,
    prompt: object,
    finished: bool,
) -> str | None

ASR text to show as the user's words, or None.

Default models do not surface Stage0. AURA uses this for a spoken turn only; vision-follow commits stay off the transcript.

validate_client_extra_body abstractmethod

validate_client_extra_body(extra_body: object) -> None

DuplexModelSessionState

Bases: ABC

Model-owned per-session state; owned by the session runner (one per session).

audio_buffer instance-attribute

audio_buffer: PcmAppendBuffer

committed_audio_operation_id instance-attribute

committed_audio_operation_id: str | None

committed_audio_payload instance-attribute

committed_audio_payload: dict[str, object] | None

committed_audio_reserved_bytes instance-attribute

committed_audio_reserved_bytes: int

context_locked instance-attribute

context_locked: bool

continuation_owner_id instance-attribute

continuation_owner_id: str | None

continuation_units instance-attribute

continuation_units: int

deferred_precreate_response instance-attribute

deferred_precreate_response: bool

deferred_response_create instance-attribute

deferred_response_create: bool

input_since_commit instance-attribute

input_since_commit: bool

last_native_submit_monotonic instance-attribute

last_native_submit_monotonic: float | None

pending_silence_owner_id instance-attribute

pending_silence_owner_id: str | None

pending_silence_task instance-attribute

pending_silence_task: Task[bool] | None

silence_deadline_monotonic instance-attribute

silence_deadline_monotonic: float | None

speech_since_commit instance-attribute

speech_since_commit: bool

clear_committed_audio abstractmethod

clear_committed_audio() -> int

clear_continuation abstractmethod

clear_continuation() -> None

retain_committed_audio abstractmethod

retain_committed_audio(
    payload: dict[str, object],
    *,
    operation_id: str | None,
    reserved_bytes: int = 0,
) -> None

DuplexRuntimeConfigError

Bases: ValueError

A model plugin rejected client-visible runtime configuration.

code instance-attribute

code = code

PartialStageForward dataclass

One downstream update the orchestrator should submit.

close_only is a final update with no new sentence. output is the model-built payload; the orchestrator does not interpret its text.

queue_close_after means this chunk still has text, but Stage1 has finished and an earlier sentence is already in flight. The text must be submitted resumable. A non-resumable submit is an end sentinel (StreamingUpdate.from_request returns None) and aborts that sentence.

close_only class-attribute instance-attribute

close_only: bool = False

is_final_update instance-attribute

is_final_update: bool

output instance-attribute

output: object

queue_close_after class-attribute instance-attribute

queue_close_after: bool = False

PcmAppendBuffer

Bases: ABC

pending_byte_count abstractmethod property

pending_byte_count: int

clear abstractmethod

clear() -> None

clear_force_listen abstractmethod

clear_force_listen() -> None

flush abstractmethod

flush(*, chunk_period_ms: int) -> dict[str, object] | None

has_pending abstractmethod

has_pending() -> bool

has_reserved abstractmethod

has_reserved() -> bool

prepare_append abstractmethod

prepare_append(
    payload: dict[str, object],
    *,
    operation_id: str,
    chunk_period_ms: int,
    allow_emit: bool,
) -> PcmAppendReservation | None

prepare_commit abstractmethod

prepare_commit(
    *, operation_id: str, chunk_period_ms: int
) -> PcmAppendReservation

PcmAppendReservation

Bases: ABC

active abstractmethod property

active: bool

byte_count abstractmethod property

byte_count: int

operation_id instance-attribute

operation_id: str

payload instance-attribute

payload: dict[str, object] | None

commit abstractmethod

commit() -> None

rollback abstractmethod

rollback() -> None

coerce_int

coerce_int(value: object) -> int | None

load_duplex_plugin

load_duplex_plugin(
    path: str, encode_audio: EncodeAudio
) -> DuplexModelPlugin

payload_turn_id

payload_turn_id(payload: object) -> int | None

reject_changed_runtime_value

reject_changed_runtime_value(
    new_value: object,
    current_value: object,
    *,
    message: str,
    code: str,
    error_cls: type[
        DuplexRuntimeConfigError
    ] = DuplexRuntimeConfigError,
) -> None

validate_duplex_plugin_sampling

validate_duplex_plugin_sampling(
    plugin: DuplexModelPlugin,
    *,
    sampling_defaults: tuple[object, ...],
) -> None

Fail fast when the plugin cannot produce one sampling parameter per stage.