diff --git a/.DS_Store b/.DS_Store index fcd63ea..eb33370 100644 Binary files a/.DS_Store and b/.DS_Store differ diff --git a/.gitignore b/.gitignore index 36b13f1..b249263 100644 --- a/.gitignore +++ b/.gitignore @@ -174,3 +174,5 @@ cython_debug/ # PyPI configuration file .pypirc +banban_server/ +talkingq-url-bak/ diff --git a/talkingq-url/assets/audio/parent_online_zh.mp3 b/talkingq-url/assets/audio/parent_online_zh.mp3 new file mode 100644 index 0000000..43dda78 Binary files /dev/null and b/talkingq-url/assets/audio/parent_online_zh.mp3 differ diff --git a/talkingq-url/banban/service/device_voice_archive.py b/talkingq-url/banban/service/device_voice_archive.py index 89a1b97..9b5b17a 100644 --- a/talkingq-url/banban/service/device_voice_archive.py +++ b/talkingq-url/banban/service/device_voice_archive.py @@ -140,10 +140,10 @@ class DeviceVoiceArchiveService(DatabaseServiceBase): ), ) stored_audio = await message_audio_storage_service.upload_audio( - sender_device_id=sender_device_id, - receiver_device_id=receiver_device_id, + device_id=sender_device_id, content=archive_audio_data, content_type=prepared_audio.mime_type, + extension=prepared_audio.archive_format, ) session_logger.info( sender_device_id, diff --git a/talkingq-url/banban/service/im.py b/talkingq-url/banban/service/im.py index 72cf2da..fa69068 100644 --- a/talkingq-url/banban/service/im.py +++ b/talkingq-url/banban/service/im.py @@ -4,11 +4,11 @@ import json from collections.abc import Mapping from pathlib import Path from typing import Any - +from services.offline_audio_cache import offline_audio_cache from fastapi import HTTPException from services.database_service_base import DatabaseServiceBase from banban.service.message_audio_storage import MessageAudioStorageService, MessageAudioStorageError - +from banban.service.binding import BindingService try: from banban.dao.im import ImDAO, DeviceIdentity, ConversationMessageCreateResult from banban.schemas.im import ( @@ -19,6 +19,9 @@ try: except ModuleNotFoundError: from banban.dao.im import ImDAO, DeviceIdentity, ConversationMessageCreateResult from banban.schemas.im import ChildConversationMessageItem, DeviceMessageCreateRequest, ParentChildMessageCreateRequest +from handlers.audio_file_handler import message_audio_storage_service +from utils.logger import session_logger + PARENT_PARTICIPANT_TYPE = 1 @@ -286,6 +289,14 @@ class ImService(DatabaseServiceBase): child_id=child_id, payload=payload, ) + binding_service = BindingService() + device = await binding_service.get_current_binding(parent_user_id) + try: + audio_url = await message_audio_storage_service.get_audio_url(stored_audio.file_key) + except Exception: + session_logger.error(device.device_id, "audio", f"failed to get audio url: {stored_audio.file_key}", exc_info=True) + audio_url = stored_audio.file_key + await offline_audio_cache.add_audio_url(device.device_id, f"{audio_url}") except Exception: try: await self.audio_storage.delete_audio(stored_audio.file_key) diff --git a/talkingq-url/banban/service/message_audio_storage.py b/talkingq-url/banban/service/message_audio_storage.py index 6d66373..c340d6a 100644 --- a/talkingq-url/banban/service/message_audio_storage.py +++ b/talkingq-url/banban/service/message_audio_storage.py @@ -142,3 +142,6 @@ class MessageAudioStorageService: ContentType=content_type, EnableMD5=False, ) + + +message_audio_storage_service = MessageAudioStorageService() diff --git a/talkingq-url/handlers/audio_file_handler.py b/talkingq-url/handlers/audio_file_handler.py index 9d19932..46cdd43 100644 --- a/talkingq-url/handlers/audio_file_handler.py +++ b/talkingq-url/handlers/audio_file_handler.py @@ -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 \ No newline at end of file + raise diff --git a/talkingq-url/handlers/mqtt_handler.py b/talkingq-url/handlers/mqtt_handler.py index e1374fc..47fe8a9 100644 --- a/talkingq-url/handlers/mqtt_handler.py +++ b/talkingq-url/handlers/mqtt_handler.py @@ -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" }, }, ) diff --git a/talkingq-url/handlers/websocket_message_handler.py b/talkingq-url/handlers/websocket_message_handler.py index 4454776..d53b950 100644 --- a/talkingq-url/handlers/websocket_message_handler.py +++ b/talkingq-url/handlers/websocket_message_handler.py @@ -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": diff --git a/talkingq-url/test/minimax_tts_more.py b/talkingq-url/test/minimax_tts_more.py index f4511c3..c7a3418 100644 --- a/talkingq-url/test/minimax_tts_more.py +++ b/talkingq-url/test/minimax_tts_more.py @@ -115,11 +115,14 @@ async def test_minimax_tts(): # "have_rest":{ # "zh": "小憩一下,待会儿见。" # } - "bind_nfc_ready":{ - "zh": "现在开始刷卡绑定设备吧。" - }, - "bind_nfc_finish":{ - "zh": "卡片绑定成功。" + # "bind_nfc_ready":{ + # "zh": "现在开始刷卡绑定设备吧。" + # }, + # "bind_nfc_finish":{ + # "zh": "卡片绑定成功。" + # } + "parent_online": { + "zh": "你好,你的家长已在线请留言" } }