Files
ClimatController/app/qingping/service.py
T

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