900 lines
23 KiB
Python
900 lines
23 KiB
Python
# 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 |