diff --git a/main/xiaozhi-server/core/providers/asr/aliyun_stream.py b/main/xiaozhi-server/core/providers/asr/aliyun_stream.py index 61ba7dc9..8246da15 100644 --- a/main/xiaozhi-server/core/providers/asr/aliyun_stream.py +++ b/main/xiaozhi-server/core/providers/asr/aliyun_stream.py @@ -96,6 +96,8 @@ class ASRProvider(ASRProviderBase): self.delete_audio_file = delete_audio_file self.expire_time = None + self.task_id = uuid.uuid4().hex + # Token管理 if self.access_key_id and self.access_key_secret: self._refresh_token() @@ -169,19 +171,23 @@ class ASRProvider(ASRProviderBase): ping_timeout=None, close_timeout=5, ) - + + self.task_id = uuid.uuid4().hex + + logger.bind(tag=TAG).info(f"WebSocket连接建立成功, task_id: {self.task_id}") + self.is_processing = True self.server_ready = False # 重置服务器准备状态 self.forward_task = asyncio.create_task(self._forward_results(conn)) - + # 发送开始请求 start_request = { "header": { "namespace": "SpeechTranscriber", "name": "StartTranscription", "status": 20000000, - "message_id": ''.join(random.choices('0123456789abcdef', k=32)), - "task_id": ''.join(random.choices('0123456789abcdef', k=32)), + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "status_text": "Gateway:SUCCESS:Success.", "appkey": self.appkey }, @@ -292,7 +298,8 @@ class ASRProvider(ASRProviderBase): "namespace": "SpeechTranscriber", "name": "StopTranscription", "status": 20000000, - "message_id": ''.join(random.choices('0123456789abcdef', k=32)), + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "status_text": "Client:Stop", "appkey": self.appkey } diff --git a/main/xiaozhi-server/core/providers/tts/aliyun_stream.py b/main/xiaozhi-server/core/providers/tts/aliyun_stream.py index 9178e2d3..42b9858e 100644 --- a/main/xiaozhi-server/core/providers/tts/aliyun_stream.py +++ b/main/xiaozhi-server/core/providers/tts/aliyun_stream.py @@ -1,3 +1,4 @@ +import random import uuid import json import hmac @@ -131,7 +132,7 @@ class TTSProvider(TTSProviderBase): self.last_active_time = None # 专属tts设置 - self.message_id = "" + self.task_id = uuid.uuid4().hex # 创建Opus编码器 self.opus_encoder = opus_encoder_utils.OpusEncoderUtils( @@ -185,7 +186,8 @@ class TTSProvider(TTSProviderBase): current_time = time.time() if self.ws and current_time - self.last_active_time < 10: # 10秒内才可以复用链接进行连续对话 - logger.bind(tag=TAG).info(f"使用已有链接...") + self.task_id = uuid.uuid4().hex + logger.bind(tag=TAG).info(f"使用已有链接..., task_id: {self.task_id}") return self.ws logger.bind(tag=TAG).info("开始建立新连接...") @@ -196,7 +198,8 @@ class TTSProvider(TTSProviderBase): ping_timeout=10, close_timeout=10, ) - logger.bind(tag=TAG).info("WebSocket连接建立成功") + self.task_id = uuid.uuid4().hex + logger.bind(tag=TAG).info(f"WebSocket连接建立成功, task_id: {self.task_id}") self.last_active_time = time.time() return self.ws except Exception as e: @@ -224,18 +227,9 @@ class TTSProvider(TTSProviderBase): if message.sentence_type == SentenceType.FIRST: # 初始化参数 try: - if not getattr(self.conn, "sentence_id", None): - self.conn.sentence_id = uuid.uuid4().hex - logger.bind(tag=TAG).info( - f"自动生成新的 会话ID: {self.conn.sentence_id}" - ) - - # aliyunStream独有的参数生成 - self.message_id = str(uuid.uuid4().hex) - logger.bind(tag=TAG).info("开始启动TTS会话...") future = asyncio.run_coroutine_threadsafe( - self.start_session(self.conn.sentence_id), + self.start_session(self.task_id), loop=self.conn.loop, ) future.result() @@ -273,7 +267,7 @@ class TTSProvider(TTSProviderBase): try: logger.bind(tag=TAG).info("开始结束TTS会话...") future = asyncio.run_coroutine_threadsafe( - self.finish_session(self.conn.sentence_id), + self.finish_session(self.task_id), loop=self.conn.loop, ) future.result() @@ -296,8 +290,8 @@ class TTSProvider(TTSProviderBase): filtered_text = MarkdownCleaner.clean_markdown(text) run_request = { "header": { - "message_id": self.message_id, - "task_id": self.conn.sentence_id, + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "namespace": "FlowingSpeechSynthesizer", "name": "RunSynthesis", "appkey": self.appkey, @@ -318,8 +312,8 @@ class TTSProvider(TTSProviderBase): self.ws = None raise - async def start_session(self, session_id): - logger.bind(tag=TAG).info(f"开始会话~~{session_id}") + async def start_session(self, task_id): + logger.bind(tag=TAG).info("开始会话~~") try: # 会话开始时检测上个会话的监听状态 if ( @@ -340,8 +334,8 @@ class TTSProvider(TTSProviderBase): start_request = { "header": { - "message_id": self.message_id, - "task_id": self.conn.sentence_id, + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "namespace": "FlowingSpeechSynthesizer", "name": "StartSynthesis", "appkey": self.appkey, @@ -365,14 +359,14 @@ class TTSProvider(TTSProviderBase): await self.close() raise - async def finish_session(self, session_id): - logger.bind(tag=TAG).info(f"关闭会话~~{session_id}") + async def finish_session(self, task_id): + logger.bind(tag=TAG).info(f"关闭会话~~{task_id}") try: if self.ws: stop_request = { "header": { - "message_id": self.message_id, - "task_id": self.conn.sentence_id, + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "namespace": "FlowingSpeechSynthesizer", "name": "StopSynthesis", "appkey": self.appkey, @@ -484,8 +478,6 @@ class TTSProvider(TTSProviderBase): loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) - # 生成会话ID - session_id = uuid.uuid4().hex # 存储音频数据 audio_data = [] @@ -504,11 +496,10 @@ class TTSProvider(TTSProviderBase): ) try: # 发送StartSynthesis请求 - start_message_id = str(uuid.uuid4().hex) start_request = { "header": { - "message_id": start_message_id, - "task_id": session_id, + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "namespace": "FlowingSpeechSynthesizer", "name": "StartSynthesis", "appkey": self.appkey, @@ -549,11 +540,10 @@ class TTSProvider(TTSProviderBase): # 发送文本合成请求 filtered_text = MarkdownCleaner.clean_markdown(text) - run_message_id = str(uuid.uuid4().hex) run_request = { "header": { - "message_id": run_message_id, - "task_id": session_id, + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "namespace": "FlowingSpeechSynthesizer", "name": "RunSynthesis", "appkey": self.appkey, @@ -563,11 +553,10 @@ class TTSProvider(TTSProviderBase): await ws.send(json.dumps(run_request)) # 发送停止合成请求 - stop_message_id = str(uuid.uuid4().hex) stop_request = { "header": { - "message_id": stop_message_id, - "task_id": session_id, + "message_id": uuid.uuid4().hex, + "task_id": self.task_id, "namespace": "FlowingSpeechSynthesizer", "name": "StopSynthesis", "appkey": self.appkey,