# from datetime import datetime, timedelta, timezone # # from bleak import BleakScanner # from bleak.backends.device import BLEDevice # from bleak.backends.scanner import AdvertisementData # # from app.qingping.models import QingpingState # from app.qingping.parser import parse_cgdn1 # # # QINGPING_SERVICE_UUID = "0000fdcd-0000-1000-8000-00805f9b34fb" # # # class QingpingService: # def __init__(self, mac: str, stale_after: float = 30.0): # self._mac = mac.upper() # self._stale_after = stale_after # # self._scanner: BleakScanner | None = None # self._state: QingpingState | None = None # # self._last_seen: datetime | None = None # self._last_error: str | None = None # # self._running = False # # @property # def state(self) -> QingpingState | None: # return self._state # # @property # def running(self) -> bool: # return self._running # # @property # def last_seen(self) -> datetime | None: # return self._last_seen # # @property # def last_error(self) -> str | None: # return self._last_error # # @property # def online(self) -> bool: # if not self._running: # return False # # if self._last_seen is None: # return False # # age = datetime.now(timezone.utc) - self._last_seen # # return age <= timedelta( # seconds=self._stale_after # ) # # async def start(self): # if self._running: # return # # try: # self._scanner = BleakScanner( # self._on_advertisement, # # service_uuids=[ # # QINGPING_SERVICE_UUID, # # ], # ) # # await self._scanner.start() # # self._running = True # self._last_error = None # # except Exception as exc: # self._scanner = None # self._running = False # self._last_error = str(exc) # # raise # # async def stop(self): # scanner = self._scanner # # self._scanner = None # self._running = False # # if scanner is not None: # await scanner.stop() # # # def _on_advertisement( # # self, # # device: BLEDevice, # # advertisement: AdvertisementData, # # ): # # if device.address.upper() != self._mac: # # return # # # # data = advertisement.service_data.get( # # QINGPING_SERVICE_UUID # # ) # # # # if not data: # # return # # # # try: # # state = parse_cgdn1( # # data, # # rssi=getattr( # # advertisement, # # "rssi", # # None, # # ), # # ) # # # # if state is None: # # return # # # # self._state = state # # # # self._last_seen = datetime.now( # # timezone.utc # # ) # # # # self._last_error = None # # # # except Exception as exc: # # self._last_error = str(exc) # # def _on_advertisement( # self, # device: BLEDevice, # advertisement: AdvertisementData, # ): # data = advertisement.service_data.get( # QINGPING_SERVICE_UUID # ) # # if not data: # return # # print( # "QINGPING:", # device.address, # device.name, # data.hex(" "), # ) # # if device.address.upper() != self._mac: # print( # "MAC mismatch:", # device.address, # "!=", # self._mac, # ) # return # # try: # state = parse_cgdn1( # data, # rssi=getattr( # advertisement, # "rssi", # None, # ), # ) # # print("PARSED:", state) # # if state is None: # return # # self._state = state # self._last_seen = datetime.now( # timezone.utc # ) # self._last_error = None # # except Exception as exc: # self._last_error = str(exc) # print("Qingping parse error:", exc) # # async def __aenter__(self): # await self.start() # return self # # async def __aexit__( # self, # exc_type, # exc_val, # exc_tb, # ): # await self.stop() import asyncio import json import logging import threading import time from dataclasses import replace from datetime import datetime, timezone import paho.mqtt.client as mqtt from .models import QingpingState from app.my_dataclasses import ( WATCHDOG_INTERVAL, QINGPING_SAMPLE_TIMEOUT, QINGPING_RECOVERY_TIMEOUT, ) logger = logging.getLogger(__name__) class QingpingService: def __init__( self, host: str, port: int = 1883, mac: str = "CCB5D131BA93", ): self._host = host self._port = port self._mac = mac self._up_topic = f"qingping/{mac}/up" self._down_topic = f"qingping/{mac}/down" self._state = QingpingState() self._lock = threading.Lock() self._client: mqtt.Client | None = None self._watchdog_task: asyncio.Task | None = None self._last_sample_monotonic: float | None = None self._device_online = False self._recovery_waiting = False self._recovery_started_monotonic: float | None = None @property def state(self) -> QingpingState: with self._lock: return self._state @property def online(self) -> bool: with self._lock: return ( self._state.mqtt_connected and self._device_online ) async def start(self) -> None: client = mqtt.Client( callback_api_version=mqtt.CallbackAPIVersion.VERSION2, client_id="tioncontroller-qingping", ) client.on_connect = self._on_connect client.on_disconnect = self._on_disconnect client.on_message = self._on_message client.reconnect_delay_set( min_delay=1, max_delay=30, ) self._client = client client.connect_async( self._host, self._port, keepalive=60, ) client.loop_start() self._watchdog_task = asyncio.create_task( self._watchdog_loop() ) async def stop(self) -> None: if self._watchdog_task is not None: self._watchdog_task.cancel() try: await self._watchdog_task except asyncio.CancelledError: pass self._watchdog_task = None if self._client is not None: self._client.disconnect() self._client.loop_stop() self._client = None def status(self) -> dict: with self._lock: state = self._state online = ( state.mqtt_connected and self._device_online ) recovery_waiting = ( self._recovery_waiting ) return { **state.to_dict(), "online": online, "recovery_waiting": recovery_waiting, } def _on_connect( self, client, _userdata, _flags, reason_code, _properties=None, ) -> None: if reason_code != 0: logger.error( "Qingping MQTT connection failed: %s", reason_code, ) return logger.info( "Qingping MQTT connected" ) client.subscribe(self._up_topic) now_monotonic = time.monotonic() with self._lock: self._state = replace( self._state, mqtt_connected=True, temperature=None, humidity=None, co2=None, pm25=None, pm10=None, battery=None, sample_timestamp=None, sample_received_at=None, ) self._last_sample_monotonic = None self._device_online = False self._recovery_waiting = True self._recovery_started_monotonic = ( now_monotonic ) logger.info( "Qingping startup initialization" ) self._send_recovery() def _on_disconnect( self, _client, _userdata, _disconnect_flags, reason_code, _properties=None, ) -> None: logger.warning( "Qingping MQTT disconnected: %s", reason_code, ) with self._lock: self._state = replace( self._state, mqtt_connected=False, ) self._device_online = False self._recovery_waiting = False self._recovery_started_monotonic = None def _on_message( self, _client, _userdata, message, ) -> None: try: payload = json.loads( message.payload.decode("utf-8") ) now = datetime.now(timezone.utc) now_monotonic = time.monotonic() message_type = int( payload.get("type") ) logger.debug( "Qingping MQTT packet received: type=%s", message_type, ) # Любой пакет означает, что устройство # физически присутствует в MQTT. with self._lock: self._state = replace( self._state, last_message_at=now, ) start_recovery = False with self._lock: if ( not self._device_online and not self._recovery_waiting ): self._state = replace( self._state, temperature=None, humidity=None, co2=None, pm25=None, pm10=None, battery=None, sample_timestamp=None, sample_received_at=None, ) self._last_sample_monotonic = None self._device_online = True self._recovery_waiting = True self._recovery_started_monotonic = ( now_monotonic ) start_recovery = True if start_recovery: logger.info( "Qingping packet received while offline: " "type=%s -> starting recovery", message_type, ) self._send_recovery() # Единственный пакет, который реально # обрабатываем как данные. if message_type == 17: self._handle_sensor_data( payload, now, ) return if message_type == 13: self._handle_heartbeat( payload, ) return logger.debug( "Qingping service packet received: type=%s", message_type, ) except Exception: logger.exception( "Qingping MQTT packet processing failed" ) def _handle_heartbeat( self, payload: dict, ) -> None: wifi_info = payload.get("wifi_info") wifi_rssi = None if isinstance(wifi_info, str): parts = wifi_info.split(",") if len(parts) >= 2: try: wifi_rssi = int(parts[1]) except ValueError: logger.warning( "Qingping invalid RSSI in wifi_info: %r", wifi_info, ) firmware = payload.get("sw_version") with self._lock: self._state = replace( self._state, wifi_rssi=wifi_rssi, firmware=firmware, ) logger.debug( "Qingping heartbeat updated: " "rssi=%s firmware=%s", wifi_rssi, firmware, ) def _handle_sensor_data( self, payload: dict, received_at: datetime, ) -> None: sensor_data = payload.get("sensorData") if not isinstance(sensor_data, list): return if not sensor_data: return sample = max( sensor_data, key=self._sample_timestamp, ) sample_timestamp = self._sample_timestamp( sample ) if sample_timestamp <= 0: return current_timestamp = ( self.state.sample_timestamp ) # CGDN1 после запуска может несколько раз # присылать одну и ту же историческую точку. if ( current_timestamp is not None and sample_timestamp <= current_timestamp ): return temperature = self._value( sample, "temperature", ) humidity = self._value( sample, "humidity", ) co2 = self._value( sample, "co2", ) pm25 = self._value( sample, "pm25", ) pm10 = self._value( sample, "pm10", ) battery = self._value( sample, "battery", ) now_monotonic = time.monotonic() with self._lock: current_timestamp = ( self._state.sample_timestamp ) # CGDN1 может повторять одну историческую точку. # Такой пакет НЕ считается новым измерением. if ( current_timestamp is not None and sample_timestamp <= current_timestamp ): logger.debug( "Qingping duplicate sensor sample ignored: " "timestamp=%s current=%s", sample_timestamp, current_timestamp, ) return was_recovering = self._recovery_waiting was_offline = not self._device_online self._state = replace( self._state, temperature=temperature, humidity=humidity, co2=co2, pm25=pm25, pm10=pm10, battery=battery, sample_timestamp=sample_timestamp, sample_received_at=received_at, ) self._last_sample_monotonic = ( now_monotonic ) self._device_online = True self._recovery_waiting = False self._recovery_started_monotonic = None # if was_offline or was_recovering: # logger.info( # "Qingping sensor stream online: " # "timestamp=%s co2=%s", # sample_timestamp, # co2, # ) # else: # logger.debug( # "Qingping sensor sample accepted: " # "timestamp=%s co2=%s", # sample_timestamp, # co2, # ) async def _watchdog_loop(self) -> None: while True: await asyncio.sleep( WATCHDOG_INTERVAL ) self._watchdog_tick( time.monotonic() ) def _watchdog_tick( self, now: float, ) -> None: recovery_timeout = None sample_timeout = None invalid_recovery_state = False with self._lock: # Ждём type 17 после recovery. if self._recovery_waiting: if ( self._recovery_started_monotonic is None ): invalid_recovery_state = True self._recovery_waiting = False self._device_online = False else: elapsed = ( now - self._recovery_started_monotonic ) if ( elapsed > QINGPING_RECOVERY_TIMEOUT ): recovery_timeout = elapsed self._recovery_waiting = False self._recovery_started_monotonic = None self._device_online = False # Нормальная работа: # следим за последним НОВЫМ type 17. elif self._last_sample_monotonic is not None: elapsed = ( now - self._last_sample_monotonic ) if ( elapsed > QINGPING_SAMPLE_TIMEOUT and self._device_online ): sample_timeout = elapsed self._device_online = False if invalid_recovery_state: logger.error( "Qingping invalid recovery state: " "recovery_waiting=True but " "recovery_started_monotonic=None" ) if recovery_timeout is not None: logger.warning( "Qingping recovery timeout: " "no type 17 for %.1f sec " "-> offline", recovery_timeout, ) if sample_timeout is not None: logger.warning( "Qingping sensor stream lost: " "no type 17 for %.1f sec " "-> offline", sample_timeout, ) @staticmethod def _sample_timestamp(sample: dict) -> int: timestamp = sample.get("timestamp") if isinstance(timestamp, dict): timestamp = timestamp.get("value") try: return int(timestamp) except (TypeError, ValueError): return 0 @staticmethod def _value(sample: dict, key: str): value = sample.get(key) if isinstance(value, dict): return value.get("value") return value async def _watchdog_loop(self) -> None: while True: await asyncio.sleep( WATCHDOG_INTERVAL ) now = time.monotonic() # -------------------------------------------------- # СЦЕНАРИЙ 1 # # Recovery уже отправлен. # Ждём type 17 максимум 30 секунд. # -------------------------------------------------- if self._recovery_waiting: if self._recovery_started_monotonic is None: logger.error( "Qingping invalid recovery state: " "recovery_waiting=True but " "recovery_started_monotonic=None" ) self._recovery_waiting = False self._device_online = False continue elapsed = now - self._recovery_started_monotonic if elapsed > QINGPING_RECOVERY_TIMEOUT: logger.warning( "Qingping recovery timeout: " "no type 17 for %.1f sec " "-> offline", elapsed, ) self._recovery_waiting = False self._recovery_started_monotonic = ( None ) self._device_online = False continue # -------------------------------------------------- # СЦЕНАРИЙ 2 # # До сих пор не получили вообще ни одного # измерения type 17. # # Ничего делать не надо. # Первый любой MQTT-пакет запустит recovery. # -------------------------------------------------- if self._last_sample_monotonic is None: continue # -------------------------------------------------- # СЦЕНАРИЙ 3 # # Нормально работали, но type 17 # перестали приходить. # -------------------------------------------------- elapsed = ( now - self._last_sample_monotonic ) if ( elapsed > QINGPING_SAMPLE_TIMEOUT and self._device_online ): logger.warning( "Qingping sensor stream lost: " "no type 17 for %.1f sec " "-> offline", elapsed, ) self._device_online = False def _send_recovery(self) -> None: if self._client is None: return payload = { "type": "17", "timestamp": int(time.time()), "setting": { "report_interval": 15, "collect_interval": 15, "need_ack": 0, }, } result = self._client.publish( self._down_topic, json.dumps( payload, separators=(",", ":"), ), ) if result.rc == mqtt.MQTT_ERR_SUCCESS: logger.warning( "Qingping recovery command sent" ) else: logger.warning( "Failed to send Qingping recovery: %s", result.rc, ) def _reset_sample_session(self) -> None: logger.debug( "Qingping sensor session reset" ) with self._lock: self._state = replace( self._state, temperature=None, humidity=None, co2=None, pm25=None, pm10=None, battery=None, sample_timestamp=None, sample_received_at=None, ) self._last_sample_monotonic = None