add test-clean code merge
This commit is contained in:
@@ -10,10 +10,10 @@ from services.offline_audio_cache import offline_audio_cache
|
||||
import aiomqtt
|
||||
from banban.service.location import location_service
|
||||
from database.models import ChildLocationCurrent
|
||||
|
||||
from datetime import datetime
|
||||
from services.device_target_cache import device_target_cache
|
||||
from utils.logger import session_logger as logger
|
||||
|
||||
from services.task_manager import task_manager
|
||||
|
||||
class TalkingQMQTTService:
|
||||
_instance = None
|
||||
@@ -95,59 +95,105 @@ class TalkingQMQTTService:
|
||||
logger.error("", "", f"消息循环异常: {e}")
|
||||
self._connected = False
|
||||
|
||||
async def _schedule_persistence(self, device_id: str, label: str, coro) -> None:
|
||||
async def _runner():
|
||||
try:
|
||||
await coro
|
||||
except Exception as exc:
|
||||
logger.error(device_id, "mqtt_persistence", f"{label} failed: {exc}", exc_info=True)
|
||||
|
||||
await task_manager.create_task(
|
||||
_runner(),
|
||||
device_id=device_id,
|
||||
task_type="persistence",
|
||||
)
|
||||
|
||||
async def _handle_device_info(self, device_id: str, payload: dict):
|
||||
# status = payload.get("status")
|
||||
data = payload.get("data", {})
|
||||
# if status == "success":
|
||||
d_id = data.get("id")
|
||||
power = data.get("power")
|
||||
signal = data.get("signal")
|
||||
version = data.get("version")
|
||||
voice = data.get("voice")
|
||||
logger.info(device_id, "", f"[设备信息] 设备 {d_id} 信息: 电量={power}, 信号强度={signal}, 版本号={version}, 音量={voice}")
|
||||
# 插入到数据库
|
||||
await device_setting_service.insert_or_update(device_id=device_id, power=power, signal_strength=signal, version_str=version, volume=voice)
|
||||
# else:
|
||||
# logger.warning(device_id, "", f"[设备信息] 设备 {device_id} 查询失败: {payload}")
|
||||
await self._schedule_persistence(
|
||||
device_id,
|
||||
"device_info",
|
||||
device_setting_service.insert_or_update(
|
||||
device_id=device_id,
|
||||
power=data.get("power"),
|
||||
signal_strength=data.get("signal"),
|
||||
version_str=data.get("version"),
|
||||
volume=data.get("voice"),
|
||||
),
|
||||
)
|
||||
await self._publish(f"device/{device_id}/event_resp", {"msg_id": "000", "status": "success"})
|
||||
|
||||
async def _handle_gps_response(self, device_id: str, payload: dict):
|
||||
status = payload.get("status")
|
||||
if payload.get("status") != "success":
|
||||
logger.warning(device_id, "", f"[GPS] query failed: {payload}")
|
||||
return
|
||||
|
||||
data = payload.get("data", {})
|
||||
if status == "success":
|
||||
lat = data.get("latitude")
|
||||
lon = data.get("longitude")
|
||||
logger.info(device_id, "", f"[GPS] 设备 {device_id} 位置: 纬度={lat}, 经度={lon}")
|
||||
# 插入到数据库
|
||||
location = ChildLocationCurrent(device_id=device_id, lat=lat, lon=lon)
|
||||
await location_service.insert_or_update(device_id=device_id, location=location)
|
||||
else:
|
||||
logger.warning(device_id, "", f"[GPS] 设备 {device_id} 查询失败: {payload}")
|
||||
raw_device_time = data.get("device_time")
|
||||
parsed_device_time = None
|
||||
if isinstance(raw_device_time, str) and raw_device_time.strip():
|
||||
try:
|
||||
parsed_device_time = datetime.fromisoformat(raw_device_time.strip())
|
||||
except ValueError:
|
||||
parsed_device_time = None
|
||||
|
||||
await self._schedule_persistence(
|
||||
device_id,
|
||||
"gps",
|
||||
location_service.report_mqtt_device_location(
|
||||
device_id=device_id,
|
||||
latitude=data.get("latitude"),
|
||||
longitude=data.get("longitude"),
|
||||
coord_type=data.get("coord_type"),
|
||||
accuracy_m=data.get("accuracy_m"),
|
||||
altitude_m=data.get("altitude_m"),
|
||||
speed_mps=data.get("speed_mps"),
|
||||
heading_deg=data.get("heading_deg"),
|
||||
source=data.get("source"),
|
||||
battery_pct=data.get("battery_pct"),
|
||||
device_time=parsed_device_time,
|
||||
),
|
||||
)
|
||||
|
||||
async def _handle_volume_response(self, device_id: str, payload: dict):
|
||||
status = payload.get("status")
|
||||
data = payload.get("data", {})
|
||||
if status == "success":
|
||||
level = data.get("current_level")
|
||||
logger.info(device_id, "", f"[音量] 设备 {device_id} 当前音量: {level}")
|
||||
# 更新设备音量
|
||||
await device_setting_service.insert_or_update(device_id=device_id, volume=level)
|
||||
if payload.get("status") != "success":
|
||||
logger.warning(device_id, "", f"[volume] command failed: {payload}")
|
||||
return
|
||||
|
||||
else:
|
||||
logger.warning(device_id, "", f"[音量] 设备 {device_id} 调节失败: {payload}")
|
||||
data = payload.get("data", {})
|
||||
current_level = data.get("current_level")
|
||||
if current_level is None:
|
||||
return
|
||||
await self._schedule_persistence(
|
||||
device_id,
|
||||
"volume",
|
||||
device_setting_service.insert_or_update(
|
||||
device_id=device_id,
|
||||
power=None,
|
||||
volume=current_level,
|
||||
signal_strength=None,
|
||||
version_str=None,
|
||||
),
|
||||
)
|
||||
|
||||
async def _handle_ota_response(self, device_id: str, payload: dict):
|
||||
status = payload.get("status")
|
||||
data = payload.get("data", {})
|
||||
if status == "accepted":
|
||||
current = data.get("current_version")
|
||||
target = data.get("target_version")
|
||||
logger.info(device_id, "", f"[OTA] 设备 {device_id} 已接受升级: {current} -> {target}")
|
||||
elif status == "success":
|
||||
new_ver = data.get("new_version")
|
||||
logger.info(device_id, "", f"[OTA] 设备 {device_id} 升级完成: {new_ver}")
|
||||
# 更新设备版本号
|
||||
await device_setting_service.insert_or_update(device_id=device_id, version_str=new_ver)
|
||||
if status == "success":
|
||||
await self._schedule_persistence(
|
||||
device_id,
|
||||
"ota",
|
||||
device_setting_service.insert_or_update(
|
||||
device_id=device_id,
|
||||
power=None,
|
||||
volume=None,
|
||||
signal_strength=None,
|
||||
version_str=data.get("new_version"),
|
||||
),
|
||||
)
|
||||
elif status != "accepted":
|
||||
logger.warning(device_id, "", f"[OTA] command failed: {payload}")
|
||||
else:
|
||||
logger.warning(device_id, "", f"[OTA] 设备 {device_id} 升级异常: {payload}")
|
||||
|
||||
|
||||
Reference in New Issue
Block a user