阿里云NLS协议中message_id每次发送都必须唯一;task_id是会话id,整个请求中需要保持一致。

This commit is contained in:
Chingfeng Li
2025-10-22 18:21:34 +08:00
parent 4b2d7da4e5
commit 2beeb825cf
2 changed files with 36 additions and 40 deletions
@@ -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
}
@@ -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,