update: 未绑定设备策略优化

fix: 音频队列竞态问题
This commit is contained in:
Sakura-RanChen
2025-12-12 18:58:24 +08:00
parent 41887ef431
commit dc170edbc1
2 changed files with 38 additions and 39 deletions
+20 -12
View File
@@ -68,7 +68,8 @@ class ConnectionHandler:
self.logger = setup_logging() self.logger = setup_logging()
self.server = server # 保存server实例的引用 self.server = server # 保存server实例的引用
self.need_bind = False # 是否需要绑定设备 self.need_bind = True # 是否需要绑定设备
self.bind_completed_event = asyncio.Event()
self.bind_code = None # 绑定设备的验证码 self.bind_code = None # 绑定设备的验证码
self.last_bind_prompt_time = 0 # 上次播放绑定提示的时间戳(秒) self.last_bind_prompt_time = 0 # 上次播放绑定提示的时间戳(秒)
self.bind_prompt_interval = 60 # 绑定提示播放间隔(秒) self.bind_prompt_interval = 60 # 绑定提示播放间隔(秒)
@@ -268,14 +269,10 @@ class ConnectionHandler:
async def _route_message(self, message): async def _route_message(self, message):
"""消息路由""" """消息路由"""
if isinstance(message, str): try:
await handleTextMessage(self, message) await asyncio.wait_for(self.bind_completed_event.wait(), timeout=1)
elif isinstance(message, bytes): except asyncio.TimeoutError:
if self.vad is None or self.asr is None: # 未绑定设备直接丢弃所有消息
return
# 未绑定设备直接丢弃所有音频,不进行ASR处理
if self.need_bind:
current_time = time.time() current_time = time.time()
# 检查是否需要播放绑定提示 # 检查是否需要播放绑定提示
if ( if (
@@ -290,6 +287,12 @@ class ConnectionHandler:
# 直接丢弃音频,不进行ASR处理 # 直接丢弃音频,不进行ASR处理
return return
if isinstance(message, str):
await handleTextMessage(self, message)
elif isinstance(message, bytes):
if self.vad is None or self.asr is None:
return
# 处理来自MQTT网关的音频包 # 处理来自MQTT网关的音频包
if self.conn_from_mqtt_gateway and len(message) >= 16: if self.conn_from_mqtt_gateway and len(message) >= 16:
handled = await self._process_mqtt_audio_message(message) handled = await self._process_mqtt_audio_message(message)
@@ -461,6 +464,9 @@ class ConnectionHandler:
self.logger.bind(tag=TAG).error(f"实例化组件失败: {e}") self.logger.bind(tag=TAG).error(f"实例化组件失败: {e}")
def _init_prompt_enhancement(self): def _init_prompt_enhancement(self):
if self.need_bind:
return
# 更新上下文信息 # 更新上下文信息
self.prompt_manager.update_context_info(self, self.client_ip) self.prompt_manager.update_context_info(self, self.client_ip)
enhanced_prompt = self.prompt_manager.build_enhanced_prompt( enhanced_prompt = self.prompt_manager.build_enhanced_prompt(
@@ -509,6 +515,9 @@ class ConnectionHandler:
def _initialize_voiceprint(self): def _initialize_voiceprint(self):
"""为当前连接初始化声纹识别""" """为当前连接初始化声纹识别"""
if self.need_bind:
return
try: try:
voiceprint_config = self.config.get("voiceprint", {}) voiceprint_config = self.config.get("voiceprint", {})
if voiceprint_config: if voiceprint_config:
@@ -548,15 +557,14 @@ class ConnectionHandler:
self.logger.bind(tag=TAG).info( self.logger.bind(tag=TAG).info(
f"{time.time() - begin_time} 秒,异步获取差异化配置成功: {json.dumps(filter_sensitive_info(private_config), ensure_ascii=False)}" f"{time.time() - begin_time} 秒,异步获取差异化配置成功: {json.dumps(filter_sensitive_info(private_config), ensure_ascii=False)}"
) )
self.need_bind = False
self.bind_completed_event.set()
except DeviceNotFoundException as e: except DeviceNotFoundException as e:
self.need_bind = True
private_config = {} private_config = {}
except DeviceBindException as e: except DeviceBindException as e:
self.need_bind = True
self.bind_code = e.bind_code self.bind_code = e.bind_code
private_config = {} private_config = {}
except Exception as e: except Exception as e:
self.need_bind = True
self.logger.bind(tag=TAG).error(f"异步获取差异化配置失败: {e}") self.logger.bind(tag=TAG).error(f"异步获取差异化配置失败: {e}")
private_config = {} private_config = {}
@@ -85,25 +85,12 @@ class AudioRateController:
elif item_type == "audio": elif item_type == "audio":
_, opus_packet = item _, opus_packet = item
# 循环等待直到时间到达
while True:
# 计算时间差 # 计算时间差
elapsed_ms = self._get_elapsed_ms() elapsed_ms = self._get_elapsed_ms()
output_ms = self.play_position output_ms = self.play_position
if elapsed_ms < output_ms: if elapsed_ms < output_ms:
# 还不到发送时间,计算等待时长 # 还不到发送时间,返回让出控制权(非阻塞)
wait_ms = output_ms - elapsed_ms
# 等待后继续检查(允许被中断)
try:
await asyncio.sleep(wait_ms / 1000)
except asyncio.CancelledError:
self.logger.bind(tag=TAG).debug("音频发送任务被取消")
raise
# 等待结束后重新检查时间(循环回到 while True)
else:
# 时间已到,跳出等待循环
break break
# 时间已到,从队列移除并发送 # 时间已到,从队列移除并发送
@@ -132,6 +119,10 @@ class AudioRateController:
async def _send_loop(): async def _send_loop():
try: try:
while True: while True:
# 只有当队列非空时才处理,否则等待队列事件
if not self.queue:
await self.queue_empty_event.wait()
await self.check_queue(send_audio_callback) await self.check_queue(send_audio_callback)
# 如果队列空了,短暂等待后再检查(避免 busy loop) # 如果队列空了,短暂等待后再检查(避免 busy loop)
await asyncio.sleep(0.01) await asyncio.sleep(0.01)