Merge branch 'main' into test-clean

# Conflicts:
#	talkingq-url/handlers/mqtt_handler.py
#	talkingq-url/handlers/websocket_message_handler.py
This commit is contained in:
stu2not
2026-05-15 09:55:56 +08:00
10 changed files with 85 additions and 37 deletions

View File

@@ -1,22 +1,27 @@
from banban.service.message_audio_storage import MessageAudioStorageService
from banban.service.message_audio_storage import StoredMessageAudio, message_audio_storage_service
from utils.logger import session_logger
message_audio_storage_service = MessageAudioStorageService()
# async def save_audio_file(audio_data: bytes, device_id: str) -> str:
# """Upload device audio to COS and return its object key."""
# try:
# stored = await message_audio_storage_service.upload_audio(
# device_id=device_id,
# content=audio_data,
# )
# session_logger.info(device_id, "audio", f"audio uploaded to COS: {stored.file_key}")
# return stored.file_key
# except Exception as e:
# session_logger.error(device_id, "audio", f"failed to store audio: {e}", exc_info=True)
# raise
async def upload_message_audio(
audio_data: bytes,
device_id: str,
*,
content_type: str = "audio/mpeg",
extension: str = "mp3",
) -> StoredMessageAudio:
"""Upload message audio to COS and return its storage metadata."""
try:
stored = await message_audio_storage_service.upload_audio(
device_id=device_id,
content=audio_data,
content_type=content_type,
extension=extension,
)
session_logger.info(device_id, "audio", f"audio uploaded to COS: {stored.file_key}")
return stored
except Exception as e:
session_logger.error(device_id, "audio", f"failed to store audio: {e}", exc_info=True)
raise
import os
@@ -54,4 +59,4 @@ async def save_audio_file(audio_data: bytes, device_id: str) -> str:
return filepath
except Exception as e:
session_logger.error(device_id, "audio", f"保存音频文件时出错: {e}", exc_info=True)
raise
raise

View File

@@ -344,7 +344,7 @@ class TalkingQMQTTService:
"status": "success",
"type": 0,
"params": {
"url": "http://101.35.224.118:8080/assets/audio/parent_online_zh.mp3"
"url": f"http://{settings.server_host}:{settings.server_port}/assets/audio/parent_online_zh.mp3"
},
}
await self._publish(f"device/{device_id}/event_resp", response_payload)
@@ -371,7 +371,7 @@ class TalkingQMQTTService:
"status": "success",
"type": 0,
"params": {
"url": "http://101.35.224.118:8080/assets/audio/parent_online_zh.mp3"
"url": f"http://{settings.server_host}:{settings.server_port}/assets/audio/parent_online_zh.mp3"
},
},
)

View File

@@ -4,7 +4,7 @@ import asyncio
from fastapi import WebSocket
from handlers.audio_packet_parser import parse_packet
from handlers.audio_session_handler import handle_websocket_data
from handlers.audio_file_handler import message_audio_storage_service, save_audio_file
from handlers.audio_file_handler import message_audio_storage_service, save_audio_file, upload_message_audio
from services.audio_session import audio_session_manager
from services.interrupt_handler import interrupt_handler
from services.task_manager import task_manager
@@ -19,9 +19,19 @@ from handlers.prompt_sound_handler import handle_prompt_sound_request
from handlers.session_cleanup_handler import handle_old_session_cleanup
from config import settings
from banban.service.im import im_service as im_conversation_service
from handlers.audio_file_handler import message_audio_storage_service, save_audio_file
from utils.audio_format import detect_audio_format, wrap_pcm_as_wav
from fastapi import HTTPException
def prepare_message_audio(audio_data: bytes) -> tuple[bytes, str, str, str]:
source_format = detect_audio_format(audio_data)
if source_format == "wav":
return audio_data, "audio/wav", "wav", source_format
if source_format == "mp3":
return audio_data, "audio/mpeg", "mp3", source_format
return wrap_pcm_as_wav(audio_data), "audio/wav", "wav", source_format
async def handle_websocket_messages(websocket: WebSocket, device_id: str, serial_number: str):
"""
处理WebSocket连接中的所有消息
@@ -267,16 +277,30 @@ async def process_parent_leave_message(device_id: str, audio_cache_key: str):
session_logger.info(device_id, "parent", "发给家长的留言没有缓存音频数据")
return
audio_file_key = await save_audio_file(cached_audio, device_id)
audio_url = f"http://{settings.server_host}:{settings.server_port}/{audio_file_key}"
archive_audio, media_mime_type, extension, source_format = prepare_message_audio(cached_audio)
stored_audio = await upload_message_audio(
archive_audio,
device_id,
content_type=media_mime_type,
extension=extension,
)
await im_conversation_service.create_device_parent_leave_message(
device_id=device_id,
media_file_key=audio_url,
media_mime_type="audio/mpeg",
media_size_bytes=len(cached_audio),
ext_json={"source": "device_ws_parent_leave_message"},
media_file_key=stored_audio.file_key,
media_mime_type=media_mime_type,
media_size_bytes=len(archive_audio),
ext_json={
"source": "device_ws_parent_leave_message",
"storage": "cos",
"source_format": source_format,
"archive_format": extension,
},
)
session_logger.info(
device_id,
"parent",
f"发给家长的留言已上传COS并写入家长会话: {stored_audio.file_key}",
)
session_logger.info(device_id, "parent", "发给家长的留言已写入家长会话")
websocket = await connection_manager.get_connection(device_id)
if websocket and websocket.client_state.name == "CONNECTED":