设备首页接入真实状态与音量控制
This commit is contained in:
@@ -60,4 +60,58 @@ class DeviceDAO(BaseDAO):
|
||||
),
|
||||
params,
|
||||
)
|
||||
return result.mappings().all()
|
||||
return result.mappings().all()
|
||||
|
||||
async def get_device_status(
|
||||
self,
|
||||
*,
|
||||
device_id: str,
|
||||
user_id: int,
|
||||
) -> Mapping[str, Any]:
|
||||
from fastapi import HTTPException
|
||||
|
||||
result = await self.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT
|
||||
db.device_id,
|
||||
db.child_id,
|
||||
c.child_name,
|
||||
ds.power,
|
||||
ds.volume,
|
||||
ds.`signal` AS signal_strength,
|
||||
ds.`version` AS version,
|
||||
ds.updated_at AS settings_updated_at,
|
||||
cl.coord_type,
|
||||
cl.lat,
|
||||
cl.lng,
|
||||
cl.accuracy_m,
|
||||
cl.altitude_m,
|
||||
cl.speed_mps,
|
||||
cl.heading_deg,
|
||||
cl.source,
|
||||
cl.battery_pct,
|
||||
cl.device_time,
|
||||
cl.server_time,
|
||||
cl.updated_at AS location_updated_at
|
||||
FROM device_bindings AS db
|
||||
LEFT JOIN children AS c
|
||||
ON c.child_id = db.child_id
|
||||
AND c.status = 1
|
||||
LEFT JOIN device_settings AS ds
|
||||
ON ds.device_id = db.device_id
|
||||
LEFT JOIN child_location_current AS cl
|
||||
ON cl.child_id = db.child_id
|
||||
AND cl.device_id = db.device_id
|
||||
WHERE db.device_id = :device_id
|
||||
AND db.owner_user_id = :user_id
|
||||
AND db.status = 1
|
||||
LIMIT 1
|
||||
"""
|
||||
),
|
||||
{"device_id": device_id, "user_id": user_id},
|
||||
)
|
||||
row = result.mappings().first()
|
||||
if row is None:
|
||||
raise HTTPException(status_code=404, detail="device not found")
|
||||
return row
|
||||
|
||||
@@ -3,7 +3,7 @@ from collections.abc import Mapping
|
||||
from datetime import datetime
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query, Request
|
||||
from pydantic import BaseModel
|
||||
from pydantic import BaseModel, Field
|
||||
from sqlalchemy import text
|
||||
|
||||
try:
|
||||
@@ -47,6 +47,39 @@ class DeviceMessageListResponse(BaseModel):
|
||||
next_cursor: int | None = None
|
||||
|
||||
|
||||
class DeviceStatusResponse(BaseModel):
|
||||
device_id: str
|
||||
child_id: int | None = None
|
||||
child_name: str | None = None
|
||||
power: int | None = None
|
||||
volume: int | None = None
|
||||
signal: int | None = None
|
||||
version: str | None = None
|
||||
settings_updated_at: datetime | None = None
|
||||
coord_type: str | None = None
|
||||
lat: float | None = None
|
||||
lng: float | None = None
|
||||
accuracy_m: int | None = None
|
||||
altitude_m: float | None = None
|
||||
speed_mps: float | None = None
|
||||
heading_deg: int | None = None
|
||||
source: int | None = None
|
||||
battery_pct: int | None = None
|
||||
device_time: datetime | None = None
|
||||
server_time: datetime | None = None
|
||||
location_updated_at: datetime | None = None
|
||||
|
||||
|
||||
class DeviceVolumeUpdateRequest(BaseModel):
|
||||
level: int = Field(ge=0, le=100)
|
||||
|
||||
|
||||
class DeviceVolumeUpdateResponse(BaseModel):
|
||||
device_id: str
|
||||
level: int
|
||||
msg_id: str
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
@@ -105,6 +138,31 @@ def _row_to_trajectory_item(row: Mapping, *, child_name: str | None) -> DeviceLo
|
||||
)
|
||||
|
||||
|
||||
def _row_to_device_status_response(row: Mapping) -> DeviceStatusResponse:
|
||||
return DeviceStatusResponse(
|
||||
device_id=str(row["device_id"]),
|
||||
child_id=int(row["child_id"]) if row["child_id"] is not None else None,
|
||||
child_name=row.get("child_name"),
|
||||
power=row["power"],
|
||||
volume=row["volume"],
|
||||
signal=row["signal_strength"],
|
||||
version=row["version"],
|
||||
settings_updated_at=row["settings_updated_at"],
|
||||
coord_type=row["coord_type"],
|
||||
lat=float(row["lat"]) if row["lat"] is not None else None,
|
||||
lng=float(row["lng"]) if row["lng"] is not None else None,
|
||||
accuracy_m=row["accuracy_m"],
|
||||
altitude_m=float(row["altitude_m"]) if row["altitude_m"] is not None else None,
|
||||
speed_mps=float(row["speed_mps"]) if row["speed_mps"] is not None else None,
|
||||
heading_deg=row["heading_deg"],
|
||||
source=int(row["source"]) if row["source"] is not None else None,
|
||||
battery_pct=row["battery_pct"],
|
||||
device_time=row["device_time"],
|
||||
server_time=row["server_time"],
|
||||
location_updated_at=row["location_updated_at"],
|
||||
)
|
||||
|
||||
|
||||
@router.get("/{device_id}/messages", response_model=DeviceMessageListResponse)
|
||||
async def list_device_messages(
|
||||
device_id: str,
|
||||
@@ -142,6 +200,57 @@ async def list_device_messages(
|
||||
)
|
||||
|
||||
|
||||
@router.get("/{device_id}/status", response_model=DeviceStatusResponse)
|
||||
async def get_device_status(
|
||||
device_id: str,
|
||||
request: Request,
|
||||
current_user_id: int = Depends(get_current_user_id),
|
||||
) -> DeviceStatusResponse:
|
||||
row = await device_service.get_device_status(
|
||||
device_id=device_id,
|
||||
user_id=current_user_id,
|
||||
)
|
||||
|
||||
logger.info(
|
||||
"device status fetched",
|
||||
extra={
|
||||
"event": "device_status",
|
||||
"request_id": getattr(request.state, "request_id", None),
|
||||
"user_id": current_user_id,
|
||||
"device_id": device_id,
|
||||
"child_id": row["child_id"],
|
||||
},
|
||||
)
|
||||
return _row_to_device_status_response(row)
|
||||
|
||||
|
||||
@router.post("/{device_id}/volume", response_model=DeviceVolumeUpdateResponse)
|
||||
async def set_device_volume(
|
||||
device_id: str,
|
||||
payload: DeviceVolumeUpdateRequest,
|
||||
request: Request,
|
||||
current_user_id: int = Depends(get_current_user_id),
|
||||
) -> DeviceVolumeUpdateResponse:
|
||||
msg_id = await device_service.set_device_volume(
|
||||
device_id=device_id,
|
||||
user_id=current_user_id,
|
||||
level=payload.level,
|
||||
)
|
||||
|
||||
logger.info(
|
||||
"device volume command sent",
|
||||
extra={
|
||||
"event": "device_volume_set",
|
||||
"request_id": getattr(request.state, "request_id", None),
|
||||
"user_id": current_user_id,
|
||||
"device_id": device_id,
|
||||
"level": payload.level,
|
||||
"msg_id": msg_id,
|
||||
},
|
||||
)
|
||||
return DeviceVolumeUpdateResponse(device_id=device_id, level=payload.level, msg_id=msg_id)
|
||||
|
||||
|
||||
@router.get("/{device_id}/location", response_model=DeviceLocationCurrentResponse)
|
||||
async def get_current_device_location(
|
||||
device_id: str,
|
||||
@@ -203,4 +312,4 @@ async def get_device_location_trajectory(
|
||||
total=len(rows),
|
||||
start_at=start_at,
|
||||
end_at=end_at,
|
||||
)
|
||||
)
|
||||
|
||||
@@ -2,6 +2,7 @@ from collections.abc import Mapping
|
||||
from typing import Any, List
|
||||
|
||||
from services.database_service_base import DatabaseServiceBase
|
||||
from fastapi import HTTPException
|
||||
|
||||
from banban.dao.device import DeviceDAO
|
||||
|
||||
@@ -38,6 +39,35 @@ class DeviceService(DatabaseServiceBase):
|
||||
finally:
|
||||
await db_session.close()
|
||||
|
||||
async def get_device_status(
|
||||
self,
|
||||
*,
|
||||
device_id: str,
|
||||
user_id: int,
|
||||
) -> Mapping[str, Any]:
|
||||
db_session = await self.get_session()
|
||||
try:
|
||||
dao = DeviceDAO(db_session)
|
||||
return await dao.get_device_status(device_id=device_id, user_id=user_id)
|
||||
finally:
|
||||
await db_session.close()
|
||||
|
||||
async def set_device_volume(
|
||||
self,
|
||||
*,
|
||||
device_id: str,
|
||||
user_id: int,
|
||||
level: int,
|
||||
) -> str:
|
||||
await self.ensure_device_access(device_id=device_id, user_id=user_id)
|
||||
|
||||
from handlers.mqtt_handler import TalkingQMQTTService
|
||||
|
||||
service = await TalkingQMQTTService.get_instance()
|
||||
if service is None:
|
||||
raise HTTPException(status_code=503, detail="MQTT 服务未初始化")
|
||||
return await service.send_volume_command(device_id, level)
|
||||
|
||||
|
||||
# 创建全局 DeviceService 实例
|
||||
device_service = DeviceService()
|
||||
device_service = DeviceService()
|
||||
|
||||
@@ -93,7 +93,14 @@ class DeviceSettingService(DatabaseServiceBase):
|
||||
finally:
|
||||
await db_session.close()
|
||||
|
||||
async def insert_or_update(self, device_id: str, power: int, volume: int, signal_strength: int, version_str: str) -> None:
|
||||
async def insert_or_update(
|
||||
self,
|
||||
device_id: str,
|
||||
power: Optional[int],
|
||||
volume: Optional[int],
|
||||
signal_strength: Optional[int],
|
||||
version_str: Optional[str],
|
||||
) -> None:
|
||||
try:
|
||||
current_row = await self.get_setting_by_device_id(device_id=device_id)
|
||||
if current_row:
|
||||
@@ -104,4 +111,4 @@ class DeviceSettingService(DatabaseServiceBase):
|
||||
pass
|
||||
|
||||
# 创建全局 DeviceSettingService 实例
|
||||
device_setting_service = DeviceSettingService()
|
||||
device_setting_service = DeviceSettingService()
|
||||
|
||||
@@ -256,10 +256,10 @@ class DeviceSetting(Base):
|
||||
timezone: Mapped[str] = mapped_column(String(32), server_default=text("'Asia/Shanghai'"))
|
||||
volume: Mapped[Optional[int]] = mapped_column(Integer)
|
||||
brightness: Mapped[Optional[int]] = mapped_column(Integer)
|
||||
disable_weekdays: Mapped[Optional[str]] = mapped_column(String(32))
|
||||
power: Mapped[Optional[int]] = mapped_column(Integer)
|
||||
signal_strength: Mapped[Optional[int]] = mapped_column(Integer)
|
||||
version_str: Mapped[Optional[str]] = mapped_column(String(64))
|
||||
signal: Mapped[Optional[int]] = mapped_column("signal", Integer)
|
||||
version: Mapped[Optional[str]] = mapped_column("version", String(64))
|
||||
disable_weekdays: Mapped[Optional[str]] = mapped_column(String(32))
|
||||
created_at: Mapped[Optional[datetime]] = mapped_column(DateTime, server_default=text("CURRENT_TIMESTAMP"))
|
||||
updated_at: Mapped[Optional[datetime]] = mapped_column(
|
||||
DateTime,
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import asyncio
|
||||
import json
|
||||
import time
|
||||
from datetime import datetime
|
||||
from typing import Awaitable, Callable, Dict, Optional
|
||||
|
||||
import aiomqtt
|
||||
@@ -9,10 +10,10 @@ from banban.service.binding import BindingService
|
||||
from banban.service.device_setting import device_setting_service
|
||||
from banban.service.location import location_service
|
||||
from config import settings
|
||||
from database.models import ChildLocationCurrent
|
||||
from services.card_service import card_service
|
||||
from services.device_target_cache import device_target_cache
|
||||
from services.offline_audio_cache import offline_audio_cache
|
||||
from services.task_manager import task_manager
|
||||
from utils.logger import session_logger as logger
|
||||
|
||||
|
||||
@@ -90,14 +91,31 @@ class TalkingQMQTTService:
|
||||
logger.error("", "", f"mqtt loop failed: {exc}")
|
||||
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):
|
||||
data = payload.get("data", {})
|
||||
await 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._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"})
|
||||
|
||||
@@ -107,12 +125,31 @@ class TalkingQMQTTService:
|
||||
return
|
||||
|
||||
data = payload.get("data", {})
|
||||
location = ChildLocationCurrent(
|
||||
device_id=device_id,
|
||||
lat=data.get("latitude"),
|
||||
lon=data.get("longitude"),
|
||||
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,
|
||||
),
|
||||
)
|
||||
await location_service.insert_or_update(device_id=device_id, location=location)
|
||||
|
||||
async def _handle_volume_response(self, device_id: str, payload: dict):
|
||||
if payload.get("status") != "success":
|
||||
@@ -120,13 +157,36 @@ class TalkingQMQTTService:
|
||||
return
|
||||
|
||||
data = payload.get("data", {})
|
||||
await device_setting_service.insert_or_update(device_id=device_id, volume=data.get("current_level"))
|
||||
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 == "success":
|
||||
await device_setting_service.insert_or_update(device_id=device_id, version_str=data.get("new_version"))
|
||||
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}")
|
||||
|
||||
|
||||
Reference in New Issue
Block a user