Files
banban/talkingq-url/services/device_config.py
2026-03-24 15:04:36 +08:00

207 lines
8.4 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import asyncio
from typing import Dict, Optional
import os
import time
from sqlalchemy import select, update, insert
from sqlalchemy.ext.asyncio import AsyncSession
from utils.logger import session_logger
from pydantic import BaseModel
from database.models import DeviceConfig as DBDeviceConfig
from services.database_service_base import DatabaseServiceBase
class DeviceConfig(BaseModel):
selected_role_key: str
preferred_language: Optional[str] = None
last_update_time: float = None
def __init__(self, **data):
if 'last_update_time' not in data:
data['last_update_time'] = time.time()
super().__init__(**data)
def to_dict(self):
"""将配置转换为可序列化的字典"""
return {
"selected_role_key": self.selected_role_key,
"preferred_language": self.preferred_language,
}
@classmethod
def from_db_model(cls, db_model: DBDeviceConfig):
"""从数据库模型创建配置对象"""
return cls(
selected_role_key=db_model.selected_role_key,
preferred_language=db_model.preferred_language,
last_update_time=db_model.last_update_time
)
class DeviceConfigManager(DatabaseServiceBase):
def __init__(self):
super().__init__(service_name="device_config")
self.device_configs: Dict[str, DeviceConfig] = {}
self.lock = asyncio.Lock()
self.worker_id = os.environ.get("UVICORN_WID", "0")
self.config_last_updated = {}
async def _is_config_updated(self, device_id: str, async_session: AsyncSession) -> bool:
"""检查配置是否已在数据库中更新"""
try:
query = select(DBDeviceConfig.last_update_time).where(DBDeviceConfig.device_id == device_id)
result = await async_session.execute(query)
db_last_updated = result.scalar_one_or_none()
if db_last_updated is None:
return False
local_last_updated = self.config_last_updated.get(device_id, 0)
return db_last_updated > local_last_updated
except Exception as e:
session_logger.error(device_id, "config", f"检查配置更新失败: {str(e)}")
return False
async def _load_config_from_db(self, device_id: str, async_session: AsyncSession):
"""从数据库加载设备配置"""
try:
query = select(DBDeviceConfig).where(DBDeviceConfig.device_id == device_id)
result = await async_session.execute(query)
db_config = result.scalar_one_or_none()
if db_config:
device_config = DeviceConfig.from_db_model(db_config)
self.config_last_updated[device_id] = db_config.last_update_time
async with self.lock:
self.device_configs[device_id] = device_config
session_logger.info(
device_id,
"config",
f"Worker {self.worker_id}: 从数据库加载设备配置,角色: {device_config.selected_role_key}, 语言: {device_config.preferred_language or '未设置'}"
)
return device_config
return None
except Exception as e:
session_logger.error(device_id, "config", f"从数据库加载配置失败: {str(e)}")
return None
async def _save_config_to_db(self, device_id: str, config: DeviceConfig, async_session: AsyncSession):
"""保存设备配置到数据库"""
try:
query = select(DBDeviceConfig).where(DBDeviceConfig.device_id == device_id)
result = await async_session.execute(query)
existing_config = result.scalar_one_or_none()
current_time = time.time()
if existing_config:
stmt = update(DBDeviceConfig).where(
DBDeviceConfig.device_id == device_id
).values(
selected_role_key=config.selected_role_key,
preferred_language=config.preferred_language,
last_update_time=current_time
)
else:
stmt = insert(DBDeviceConfig).values(
device_id=device_id,
selected_role_key=config.selected_role_key,
preferred_language=config.preferred_language,
last_update_time=current_time
)
await async_session.execute(stmt)
await async_session.commit()
self.config_last_updated[device_id] = current_time
session_logger.info(
device_id,
"config",
f"Worker {self.worker_id}: 设备配置已保存到数据库"
)
except Exception as e:
await async_session.rollback()
session_logger.error(
device_id,
"config",
f"Worker {self.worker_id}: 保存设备配置到数据库失败: {str(e)}"
)
raise
async def get_config(self, device_id: str, force_refresh: bool = False) -> DeviceConfig:
"""获取设备配置如果force_refresh=True则强制从数据库读取"""
await self._init_database()
config = None
config_updated = False
db_session = await self.db_manager.get_session()
try:
if force_refresh or await self._is_config_updated(device_id, db_session):
config = await self._load_config_from_db(device_id, db_session)
config_updated = True
if not config_updated:
async with self.lock:
config = self.device_configs.get(device_id)
if not config:
from config import settings
default_language = "en"
config = DeviceConfig(
selected_role_key=settings.selected_role_key,
preferred_language=default_language
)
session_logger.info(
device_id,
"config",
f"Worker {self.worker_id}: 为新设备创建配置,默认角色: {settings.selected_role_key},默认语言: {default_language}"
)
await self.set_config(device_id, config)
elif config_updated:
session_logger.info(
device_id,
"config",
f"Worker {self.worker_id}: 已从数据库刷新设备配置,当前角色: {config.selected_role_key}, 语言: {config.preferred_language or '未设置'}"
)
return config
except Exception as e:
session_logger.error(device_id, "config", f"获取设备配置失败: {str(e)}")
raise
finally:
await db_session.close()
async def set_config(self, device_id: str, config: DeviceConfig):
"""设置并保存设备配置"""
await self._init_database()
config.last_update_time = time.time()
db_session = await self.db_manager.get_session()
try:
async with self.lock:
old_config = self.device_configs.get(device_id)
self.device_configs[device_id] = config
if old_config:
changes = []
if old_config.selected_role_key != config.selected_role_key:
changes.append(f"角色: {old_config.selected_role_key} -> {config.selected_role_key}")
if old_config.preferred_language != config.preferred_language:
changes.append(f"语言: {old_config.preferred_language or '未设置'} -> {config.preferred_language or '未设置'}")
if changes:
session_logger.info(device_id, "config", f"设备配置已更新: {', '.join(changes)}")
await self._save_config_to_db(device_id, config, db_session)
except Exception as e:
session_logger.error(device_id, "config", f"设置设备配置失败: {str(e)}")
raise
finally:
await db_session.close()
device_config_manager = DeviceConfigManager()