diff --git a/main/xiaozhi-server/core/handle/sendAudioHandle.py b/main/xiaozhi-server/core/handle/sendAudioHandle.py index d0167498..8198dfb0 100644 --- a/main/xiaozhi-server/core/handle/sendAudioHandle.py +++ b/main/xiaozhi-server/core/handle/sendAudioHandle.py @@ -7,6 +7,10 @@ from core.providers.tts.dto.dto import SentenceType from core.utils.audioRateController import AudioRateController TAG = __name__ +# 音频帧时长(毫秒) +AUDIO_FRAME_DURATION = 60 +# 预缓冲包数量,直接发送以减少延迟 +PRE_BUFFER_COUNT = 5 async def sendAudioMessage(conn, sentenceType, audios, text): @@ -45,7 +49,7 @@ async def sendAudioMessage(conn, sentenceType, audios, text): async def _wait_for_audio_completion(conn): """ - 等待音频队列清空 + 等待音频队列清空并等待预缓冲包播放完成 Args: conn: 连接对象 @@ -56,6 +60,13 @@ async def _wait_for_audio_completion(conn): f"等待音频发送完成,队列中还有 {len(rate_controller.queue)} 个包" ) await rate_controller.queue_empty_event.wait() + + # 等待预缓冲包播放完成 + # 前N个包直接发送,增加2个网络抖动包,需要额外等待它们在客户端播放完成 + frame_duration_ms = rate_controller.frame_duration + pre_buffer_playback_time = (PRE_BUFFER_COUNT + 2) * frame_duration_ms / 1000.0 + await asyncio.sleep(pre_buffer_playback_time) + conn.logger.bind(tag=TAG).debug("音频发送完成") @@ -81,14 +92,14 @@ async def _send_to_mqtt_gateway(conn, opus_packet, timestamp, sequence): await conn.websocket.send(complete_packet) -async def sendAudio(conn, audios, frame_duration=60): +async def sendAudio(conn, audios, frame_duration=AUDIO_FRAME_DURATION): """ 发送音频包,使用 AudioRateController 进行精确的流量控制 Args: conn: 连接对象 audios: 单个opus包(bytes) 或 opus包列表 - frame_duration: 帧时长(毫秒),默认60ms + frame_duration: 帧时长(毫秒),默认使用全局常量AUDIO_FRAME_DURATION """ if audios is None or len(audios) == 0: return @@ -187,16 +198,14 @@ async def _send_audio_with_rate_control( flow_control: 流控状态 send_delay: 固定延迟(秒),-1表示使用动态流控 """ - pre_buffer_count = 5 - for packet in audio_list: if conn.client_abort: return conn.last_activity_time = time.time() * 1000 - # 预缓冲:前5个包直接发送 - if flow_control["packet_count"] < pre_buffer_count: + # 预缓冲:前N个包直接发送 + if flow_control["packet_count"] < PRE_BUFFER_COUNT: await _do_send_audio(conn, packet, flow_control) conn.client_is_speaking = True elif send_delay > 0: diff --git a/main/xiaozhi-server/core/utils/audioRateController.py b/main/xiaozhi-server/core/utils/audioRateController.py index 35038d27..01d71b05 100644 --- a/main/xiaozhi-server/core/utils/audioRateController.py +++ b/main/xiaozhi-server/core/utils/audioRateController.py @@ -1,5 +1,6 @@ import time import asyncio +from collections import deque from config.logger import setup_logging TAG = __name__ @@ -18,13 +19,14 @@ class AudioRateController: frame_duration: 单个音频帧时长(毫秒),默认60ms """ self.frame_duration = frame_duration - self.queue = [] + self.queue = deque() self.play_position = 0 # 虚拟播放位置(毫秒) self.start_timestamp = None # 开始时间戳(只读,不修改) self.pending_send_task = None self.logger = logger self.queue_empty_event = asyncio.Event() # 队列清空事件 self.queue_empty_event.set() # 初始为空状态 + self.queue_has_data_event = asyncio.Event() # 队列数据事件 def reset(self): """重置控制器状态""" @@ -34,13 +36,17 @@ class AudioRateController: self.queue.clear() self.play_position = 0 - self.start_timestamp = time.time() - self.queue_empty_event.set() # 队列已清空 + self.start_timestamp = None # 由首个音频包设置 + # 相关事件处理 + self.queue_empty_event.set() + self.queue_has_data_event.clear() def add_audio(self, opus_packet): """添加音频包到队列""" self.queue.append(("audio", opus_packet)) - self.queue_empty_event.clear() # 队列非空,清除事件 + # 相关事件处理 + self.queue_empty_event.clear() + self.queue_has_data_event.set() def add_message(self, message_callback): """ @@ -50,13 +56,15 @@ class AudioRateController: message_callback: 消息发送回调函数 async def() """ self.queue.append(("message", message_callback)) - self.queue_empty_event.clear() # 队列非空,清除事件 + # 相关事件处理 + self.queue_empty_event.clear() + self.queue_has_data_event.set() def _get_elapsed_ms(self): """获取已经过的时间(毫秒)""" if self.start_timestamp is None: return 0 - return (time.time() - self.start_timestamp) * 1000 + return (time.monotonic() - self.start_timestamp) * 1000 async def check_queue(self, send_audio_callback): """ @@ -65,9 +73,6 @@ class AudioRateController: Args: send_audio_callback: 发送音频的回调函数 async def(opus_packet) """ - if self.start_timestamp is None: - self.start_timestamp = time.time() - while self.queue: item = self.queue[0] item_type = item[0] @@ -75,7 +80,7 @@ class AudioRateController: if item_type == "message": # 消息类型:立即发送,不占用播放时间 _, message_callback = item - self.queue.pop(0) + self.queue.popleft() try: await message_callback() except Exception as e: @@ -83,6 +88,9 @@ class AudioRateController: raise elif item_type == "audio": + if self.start_timestamp is None: + self.start_timestamp = time.monotonic() + _, opus_packet = item # 循环等待直到时间到达 @@ -107,16 +115,17 @@ class AudioRateController: break # 时间已到,从队列移除并发送 - self.queue.pop(0) + self.queue.popleft() self.play_position += self.frame_duration - try: await send_audio_callback(opus_packet) except Exception as e: self.logger.bind(tag=TAG).error(f"发送音频失败: {e}") raise + # 队列处理完后清除事件 self.queue_empty_event.set() + self.queue_has_data_event.clear() def start_sending(self, send_audio_callback): """ @@ -132,9 +141,10 @@ class AudioRateController: async def _send_loop(): try: while True: + # 等待队列数据事件,不轮询等待占用CPU + await self.queue_has_data_event.wait() + await self.check_queue(send_audio_callback) - # 如果队列空了,短暂等待后再检查(避免 busy loop) - await asyncio.sleep(0.01) except asyncio.CancelledError: self.logger.bind(tag=TAG).debug("音频发送循环已停止") except Exception as e: