11from __future__ import annotations
22
3+ import time
34from typing import Generator , Iterator , TYPE_CHECKING
4- from queue import Queue , Empty , Full
5+ from queue import Queue , Empty
56from contextlib import contextmanager
7+ from threading import Lock
8+ from collections import deque
69
710import numpy as np
811import numpy .typing as npt
1720class AudioOut :
1821 output_rate : float = 48000.0 # Hz
1922 buffersize : float = 0.05 # seconds
23+ pi_kp : float = 0.002
24+ pi_ki : float = 0.003
25+ pi_deadband : float = 0.01
26+ pi_i_clamp : float = 0.005
2027
2128 input_rate : float
2229 speed : float
@@ -30,12 +37,49 @@ def __init__(
3037 self .input_rate = input_rate
3138 self .speed = speed
3239 self .resampler = resampler
33- self .queue = Queue (maxsize = 100 )
40+ self .queue = Queue ()
41+ self .queued_frames = 0
3442 self .offset = 0
43+ self .lock = Lock ()
44+
45+ # Record last time `self.send()` was called
46+ self .last_send = time .perf_counter ()
47+
48+ # Calculate initial ratio and target frames in buffer
49+ self .base_ratio = self .output_rate / self .input_rate / self .speed
50+ self .pi_ratio = self .base_ratio
51+ self .pi_target = int (round (self .output_rate * self .buffersize * 2 ))
52+ self .pi_integral = 0.0
53+ self .pi_deque = deque [int ](maxlen = 3 )
54+
55+ def adapt_ratio (self , extra_frames : int ) -> None :
56+ # Get consistent view of time and queued frames
57+ with self .lock :
58+ delta = time .perf_counter () - self .last_send
59+ queued_frames = self .queued_frames
60+
61+ # Compute pending frames
62+ frames_predicted = int (round (delta * self .output_rate ))
63+ pending_frames = extra_frames + queued_frames + frames_predicted
64+
65+ # Average pending frames over last few calls to smooth out noise
66+ self .pi_deque .appendleft (pending_frames )
67+ pending_frames = sum (self .pi_deque ) // len (self .pi_deque )
68+
69+ # Compute PI correction
70+ error = self .pi_target - pending_frames
71+ normalized_error = error / self .pi_target
72+ if abs (normalized_error ) >= self .pi_deadband :
73+ self .pi_integral += normalized_error
74+ p_term = self .pi_kp * normalized_error
75+ i_term = np .clip (
76+ self .pi_ki * self .pi_integral , - self .pi_i_clamp , self .pi_i_clamp
77+ )
78+ self .pi_integral = i_term / self .pi_ki
79+ correction = p_term + i_term
3580
36- @property
37- def ratio (self ) -> float :
38- return self .output_rate / self .input_rate / self .speed
81+ # Update ratio
82+ self .pi_ratio = self .base_ratio * (1.0 + correction )
3983
4084 @contextmanager
4185 def run (self ) -> Iterator [AudioOut ]:
@@ -51,33 +95,40 @@ def run(self) -> Iterator[AudioOut]:
5195 buffersize_msec = int (round (self .buffersize * 1000 )),
5296 )
5397 device .start (stream )
54- try :
98+ with device :
5599 yield self
56- finally :
57- device .stop ()
58100
59101 def send (self , audio : npt .NDArray [np .int16 ]) -> None :
60102 # Resample to output rate
61- data = self .resampler .process (audio , self .ratio )
103+ data = self .resampler .process (audio , self .pi_ratio )
62104 # Reshape
63105 data = data .astype (np .int16 ).reshape (- 1 , 2 )
64106 # Reduce volume
65107 data //= 4
66- # Send without blocking if possible
67- try :
108+ # Send without blocking
109+ with self . lock :
68110 self .queue .put_nowait (data )
69- # Synchronization issue, let it regulate itself
70- except Full :
71- pass
111+ self .queued_frames += len (data )
112+ self .last_send = time .perf_counter ()
72113
73114 def _audio_stream (self ) -> Generator [bytes , int , None ]:
74- # Get first required frames
115+ # Initialize state
75116 extra_frames_start = 0
76117 buffer = np .empty ((0 , 2 ), dtype = np .int16 )
118+
119+ # Get first required frames
77120 required_frames = yield b""
121+ result = np .zeros ((required_frames , 2 ), dtype = np .int16 )
122+
123+ # Wait until we have enough frames to fill the first request
124+ while self .queued_frames < 1.5 * required_frames :
125+ required_frames = yield result .tobytes ()
78126
79127 # Loop over requested frames
80128 while True :
129+ # Adapt ratio
130+ self .adapt_ratio (len (buffer ) - extra_frames_start )
131+
81132 # Initialize result buffer
82133 result = np .zeros ((required_frames , 2 ), dtype = np .int16 )
83134
@@ -96,7 +147,9 @@ def _audio_stream(self) -> Generator[bytes, int, None]:
96147 # Get more frames until we have enough
97148 while True :
98149 try :
99- buffer = self .queue .get_nowait ()
150+ with self .lock :
151+ buffer = self .queue .get_nowait ()
152+ self .queued_frames -= len (buffer )
100153 except Empty :
101154 extra_frames_start = len (buffer )
102155 break
0 commit comments