Skip to content

vllm_omni.worker_v2.streaming_audio

Request-owned PCM accumulation for a stage that decodes audio in place.

AudioRequest dataclass

active class-attribute instance-attribute

active: bool = True

chunk_frames instance-attribute

chunk_frames: int

direct_first_pending class-attribute instance-attribute

direct_first_pending: bool = False

emitted class-attribute instance-attribute

emitted: bool = False

first_audio_required class-attribute instance-attribute

first_audio_required: bool = False

lock class-attribute instance-attribute

lock: Lock = field(default_factory=threading.Lock)

pending class-attribute instance-attribute

pending: list[Tensor | ndarray] = field(
    default_factory=list
)

pending_samples class-attribute instance-attribute

pending_samples: int = 0

samples_per_frame instance-attribute

samples_per_frame: int

accept_first_audio

accept_first_audio() -> None

finish

finish() -> None

push

push(pcm: Tensor | ndarray, ended: bool) -> Tensor | None

StreamingAudioBuffer

chunk_frames instance-attribute

chunk_frames = chunk_frames

requests instance-attribute

requests: dict[str, AudioRequest] = {}

samples_per_frame instance-attribute

samples_per_frame = samples_per_frame

add

add(request_id: str) -> AudioRequest

finish

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

StreamingAudioOutput dataclass

event instance-attribute

event: Event | None

length_end instance-attribute

length_end: ndarray

query_start_loc instance-attribute

query_start_loc: ndarray

requests instance-attribute

requests: list[AudioRequest | None]

sample_rate instance-attribute

sample_rate: Tensor

valid instance-attribute

valid: Tensor

wav instance-attribute

wav: Tensor

get_output

get_output() -> list[dict[str, Tensor] | None]

to_cpu

to_cpu(
    copy_stream: Stream,
    copy_tensor: Callable[[Tensor], Tensor],
) -> StreamingAudioOutput