Железный Егорин — форк Орбиты: тёплый стиль (amber), логин с odyssey WebGL-молнией справа, энергомониторинг+панель
Ребренд Мониторинг Орбита -> Железный Егорин; accent green->amber #F59E0B; LoginScreen правая половина = Lightning (WebGL, hue 35 тёплый) под стеклянной формой; левая панель — тёплый градиент. Доступ агентам: репо публичный, push german-токеном. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Binary file not shown.
@@ -0,0 +1,51 @@
|
||||
# Приёмный (ingest) стек + шлюз на vm-mts1. Лёгкий: без TimescaleDB/HA/Grafana.
|
||||
# Устройства (трекеры Navtelecom) шлют FLEX на :2000 → flex-server декодирует →
|
||||
# MQTT (mosquitto) → api (шлюз) → REST/WS для фронта.
|
||||
#
|
||||
# Раскладка каталога (собирает деплой-агент на ноде, /opt/monitor-ingest):
|
||||
# docker-compose.yml
|
||||
# mosquitto/mosquitto.conf
|
||||
# flex-server/ (содержимое navtelecom-flex-server: Dockerfile, flex_server/, requirements.txt)
|
||||
# api/ (содержимое servaki-monitor/server: Dockerfile, package.json, src/)
|
||||
services:
|
||||
mosquitto:
|
||||
image: eclipse-mosquitto:2
|
||||
container_name: monitor-mosquitto
|
||||
restart: unless-stopped
|
||||
volumes:
|
||||
- ./mosquitto/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro
|
||||
|
||||
flex-server:
|
||||
build: ./flex-server
|
||||
container_name: monitor-flex
|
||||
restart: unless-stopped
|
||||
ports:
|
||||
- "2000:2000" # ← порт приёма устройств (указывается в трекере)
|
||||
environment:
|
||||
FLEX_LISTEN_PORT: 2000
|
||||
MQTT_ENABLED: "true"
|
||||
MQTT_HOST: mosquitto
|
||||
MQTT_PORT: 1883
|
||||
MQTT_BASE_TOPIC: navtelecom
|
||||
DB_ENABLED: "false" # без TimescaleDB — экономим диск
|
||||
LOG_LEVEL: INFO
|
||||
depends_on:
|
||||
- mosquitto
|
||||
|
||||
api:
|
||||
build: ./api
|
||||
container_name: monitor-api
|
||||
restart: unless-stopped
|
||||
ports:
|
||||
- "127.0.0.1:4100:4100" # наружу только через Caddy (/api, /socket)
|
||||
env_file:
|
||||
- .env # JWT_SECRET + YANDEX_* (секреты, не в репозитории)
|
||||
environment:
|
||||
PORT: 4100
|
||||
DATA_DIR: /data
|
||||
MQTT_URL: mqtt://mosquitto:1883
|
||||
MQTT_BASE_TOPIC: navtelecom
|
||||
volumes:
|
||||
- ./data:/data # персист арендаторов/привязок устройств
|
||||
depends_on:
|
||||
- mosquitto
|
||||
@@ -0,0 +1,3 @@
|
||||
listener 1883
|
||||
allow_anonymous true
|
||||
persistence false
|
||||
@@ -0,0 +1,9 @@
|
||||
[Unit]
|
||||
Description=Ingress link watchdog (MQTT broker reachability alert)
|
||||
After=network-online.target
|
||||
Wants=network-online.target
|
||||
|
||||
[Service]
|
||||
Type=oneshot
|
||||
ExecStart=/opt/ingress-watchdog/watchdog.sh
|
||||
Nice=10
|
||||
@@ -0,0 +1,153 @@
|
||||
#!/usr/bin/env bash
|
||||
# ingress-watchdog.sh — алерт обрыва связи ингриса с MQTT-брокером (Milestone B)
|
||||
# stdlib/bash only. Аддитивно. НЕ трогает приёмники/буферизацию/маршруты.
|
||||
#
|
||||
# Проба: TCP-connect к брокеру vm-mts1 10.50.0.48:1883 (таймаут 5с).
|
||||
# Дебаунс: DOWN только после непрерывной недоступности > DOWN_THRESHOLD (60с).
|
||||
# Алерт ТОЛЬКО на переходе состояния (up->down = ОБРЫВ, down->up = ВОССТАНОВЛЕНО).
|
||||
# Каналы: всегда лог + статус-файл; Telegram — если в .env есть TG_TOKEN и TG_CHAT_ID.
|
||||
|
||||
set -u
|
||||
|
||||
# --- Конфиг (можно переопределить через .env) ---
|
||||
BROKER_HOST="10.50.0.48"
|
||||
BROKER_PORT="1883"
|
||||
DOWN_THRESHOLD=60 # сек непрерывной недоступности до алерта ОБРЫВ
|
||||
PROBE_TIMEOUT=5 # сек на TCP-пробу
|
||||
STATE_FILE="/var/run/ingress-watchdog.state"
|
||||
STATUS_FILE="/run/ingress-watchdog.status"
|
||||
LOG_FILE="/var/log/ingress-watchdog.log"
|
||||
ENV_FILE="/opt/ingress-watchdog/.env"
|
||||
NODE_NAME="$(hostname)"
|
||||
SPOOL_FILE="" # путь к queue.jsonl (задаётся в .env на каждой ноде)
|
||||
ALERT_WEBHOOK="https://kpgen.servaki.online/webhook/alert" # единый n8n-хаб (можно переопределить в .env)
|
||||
|
||||
# --- Загрузка .env (TG_TOKEN, TG_CHAT_ID, SPOOL_FILE, NODE_NAME override) ---
|
||||
if [ -f "$ENV_FILE" ]; then
|
||||
# shellcheck disable=SC1090
|
||||
. "$ENV_FILE"
|
||||
fi
|
||||
|
||||
now="$(date +%s)"
|
||||
ts="$(date '+%Y-%m-%d %H:%M:%S %Z')"
|
||||
|
||||
log() {
|
||||
# НЕ логируем токен. Пишем строку в лог.
|
||||
printf '%s [%s] %s\n' "$ts" "$NODE_NAME" "$1" >> "$LOG_FILE" 2>/dev/null
|
||||
}
|
||||
|
||||
spool_count() {
|
||||
if [ -n "${SPOOL_FILE:-}" ] && [ -f "$SPOOL_FILE" ]; then
|
||||
wc -l < "$SPOOL_FILE" 2>/dev/null | tr -d ' '
|
||||
else
|
||||
echo 0
|
||||
fi
|
||||
}
|
||||
|
||||
send_telegram() {
|
||||
# $1 = текст. Отправляем только если заданы token+chat_id.
|
||||
local text="$1"
|
||||
if [ -n "${TG_TOKEN:-}" ] && [ -n "${TG_CHAT_ID:-}" ]; then
|
||||
timeout 10 curl -s -o /dev/null \
|
||||
"https://api.telegram.org/bot${TG_TOKEN}/sendMessage" \
|
||||
--data-urlencode "chat_id=${TG_CHAT_ID}" \
|
||||
--data-urlencode "text=${text}" >/dev/null 2>&1 \
|
||||
&& log "telegram: sent" \
|
||||
|| log "telegram: send FAILED"
|
||||
else
|
||||
log "telegram: skipped (no token/chat_id)"
|
||||
fi
|
||||
}
|
||||
|
||||
send_webhook() {
|
||||
# $1=severity (critical|warning|info), $2=title, $3=text. Единый хаб алертов (n8n).
|
||||
local sev="$1" title="$2" text="$3" iso body
|
||||
[ -z "${ALERT_WEBHOOK:-}" ] && { log "webhook: skipped (no ALERT_WEBHOOK)"; return; }
|
||||
iso="$(date -u '+%Y-%m-%dT%H:%M:%SZ')"
|
||||
body="$(printf '{"source":"orbita-ingress","node":"%s","severity":"%s","title":"%s","text":"%s","ts":"%s"}' \
|
||||
"$NODE_NAME" "$sev" "$title" "$text" "$iso")"
|
||||
local code
|
||||
code="$(timeout 12 curl -s -o /dev/null -m 10 --retry 2 -w '%{http_code}' \
|
||||
-X POST "$ALERT_WEBHOOK" -H 'Content-Type: application/json' -d "$body" 2>/dev/null)"
|
||||
if [ "$code" = "200" ]; then log "webhook: sent ($sev) HTTP 200"; else log "webhook: send FAILED (HTTP ${code:-000})"; fi
|
||||
}
|
||||
|
||||
write_status() {
|
||||
# $1=state up|down, $2=human message
|
||||
{
|
||||
printf 'node=%s\n' "$NODE_NAME"
|
||||
printf 'state=%s\n' "$1"
|
||||
printf 'updated=%s\n' "$ts"
|
||||
printf 'updated_epoch=%s\n' "$now"
|
||||
printf 'broker=%s:%s\n' "$BROKER_HOST" "$BROKER_PORT"
|
||||
printf 'spool_lines=%s\n' "$(spool_count)"
|
||||
printf 'last_message=%s\n' "$2"
|
||||
} > "$STATUS_FILE" 2>/dev/null
|
||||
}
|
||||
|
||||
# --- Читаем предыдущий стейт ---
|
||||
last_ok=0
|
||||
current_state="up"
|
||||
down_since=0
|
||||
if [ -f "$STATE_FILE" ]; then
|
||||
# shellcheck disable=SC1090
|
||||
. "$STATE_FILE" 2>/dev/null || true
|
||||
fi
|
||||
|
||||
# --- TCP-проба ---
|
||||
if timeout "$PROBE_TIMEOUT" bash -c "cat < /dev/null > /dev/tcp/${BROKER_HOST}/${BROKER_PORT}" 2>/dev/null; then
|
||||
probe_ok=1
|
||||
else
|
||||
probe_ok=0
|
||||
fi
|
||||
|
||||
new_state="$current_state"
|
||||
|
||||
if [ "$probe_ok" -eq 1 ]; then
|
||||
last_ok="$now"
|
||||
if [ "$current_state" = "down" ]; then
|
||||
# ВОССТАНОВЛЕНО (переход down->up)
|
||||
new_state="up"
|
||||
down_since=0
|
||||
spc="$(spool_count)"
|
||||
msg="✅ ВОССТАНОВЛЕНО: связь ингриса ${NODE_NAME} с брокером ${BROKER_HOST}:${BROKER_PORT} восстановлена. В спуле сейчас: ${spc} сообщ. (дренаж до-дошлёт). Время: ${ts}"
|
||||
log "STATE up (recovered). spool=${spc}"
|
||||
write_status "up" "recovered"
|
||||
send_webhook "info" "Связь восстановлена ${NODE_NAME}" "брокер ${BROKER_HOST}:${BROKER_PORT} снова доступен; в буфере ${spc} сообщ. (дренаж до-дошлёт)"
|
||||
send_telegram "$msg"
|
||||
else
|
||||
# остаёмся up, без алерта
|
||||
write_status "up" "ok"
|
||||
fi
|
||||
else
|
||||
# недоступно
|
||||
if [ "$down_since" -eq 0 ] 2>/dev/null || [ -z "${down_since:-}" ]; then
|
||||
down_since="$now"
|
||||
fi
|
||||
down_for=$(( now - down_since ))
|
||||
if [ "$current_state" = "up" ] && [ "$down_for" -gt "$DOWN_THRESHOLD" ]; then
|
||||
# ОБРЫВ (переход up->down после дебаунса)
|
||||
new_state="down"
|
||||
spc="$(spool_count)"
|
||||
msg="🔴 ОБРЫВ: ингрис ${NODE_NAME} потерял связь с брокером ${BROKER_HOST}:${BROKER_PORT} (недоступен ${down_for}с). В спуле накоплено: ${spc} сообщ. и буферизуется. Время: ${ts}"
|
||||
log "STATE down (link lost ${down_for}s). spool=${spc}"
|
||||
write_status "down" "link_lost"
|
||||
send_webhook "critical" "Обрыв связи ${NODE_NAME}" "брокер ${BROKER_HOST}:${BROKER_PORT} недоступен ${down_for}с; в буфере ${spc} сообщ., буферизуется"
|
||||
send_telegram "$msg"
|
||||
elif [ "$current_state" = "down" ]; then
|
||||
# уже в down — не спамим
|
||||
write_status "down" "still_down"
|
||||
else
|
||||
# недоступно, но порог ещё не перевален — ждём
|
||||
write_status "up" "probing_down ${down_for}s"
|
||||
fi
|
||||
fi
|
||||
|
||||
# --- Сохраняем стейт ---
|
||||
{
|
||||
printf 'last_ok=%s\n' "$last_ok"
|
||||
printf 'current_state=%s\n' "$new_state"
|
||||
printf 'down_since=%s\n' "$down_since"
|
||||
} > "$STATE_FILE" 2>/dev/null
|
||||
|
||||
exit 0
|
||||
@@ -0,0 +1,11 @@
|
||||
[Unit]
|
||||
Description=Run ingress link watchdog every 30s
|
||||
|
||||
[Timer]
|
||||
OnBootSec=45
|
||||
OnUnitActiveSec=30s
|
||||
AccuracySec=5s
|
||||
Unit=ingress-watchdog.service
|
||||
|
||||
[Install]
|
||||
WantedBy=timers.target
|
||||
@@ -0,0 +1,39 @@
|
||||
#!/usr/bin/env python3
|
||||
# Мини raw-MQTT subscriber: коннект к брокеру, SUBSCRIBE guard/#, ждём 1 PUBLISH.
|
||||
import socket, struct, sys, time
|
||||
|
||||
HOST, PORT = "10.50.0.48", 1883
|
||||
TOPIC = "guard/#"
|
||||
|
||||
def rl(n):
|
||||
o=bytearray()
|
||||
while True:
|
||||
b=n%128; n//=128
|
||||
if n>0: b|=0x80
|
||||
o.append(b)
|
||||
if n==0: break
|
||||
return bytes(o)
|
||||
def s(x):
|
||||
b=x.encode(); return struct.pack("!H",len(b))+b
|
||||
|
||||
sock=socket.create_connection((HOST,PORT),timeout=8)
|
||||
# CONNECT
|
||||
body=s("MQTT")+bytes([0x04,0x02])+struct.pack("!H",60)+s("pcn-subtest")
|
||||
sock.sendall(bytes([0x10])+rl(len(body))+body)
|
||||
ca=sock.recv(4)
|
||||
assert ca[0]==0x20 and ca[3]==0x00, ("CONNACK",ca)
|
||||
# SUBSCRIBE (packet id 1, qos0)
|
||||
sb=struct.pack("!H",1)+s(TOPIC)+bytes([0x00])
|
||||
sock.sendall(bytes([0x82])+rl(len(sb))+sb)
|
||||
# SUBACK
|
||||
sock.recv(5)
|
||||
print("SUBSCRIBED", flush=True)
|
||||
# ждём PUBLISH до 12с
|
||||
sock.settimeout(12)
|
||||
buf=b""
|
||||
try:
|
||||
data=sock.recv(4096)
|
||||
print("GOT_FROM_BROKER:", repr(data[:300]))
|
||||
except socket.timeout:
|
||||
print("TIMEOUT_NO_MESSAGE")
|
||||
sock.close()
|
||||
@@ -0,0 +1,9 @@
|
||||
// Вписывает settings.errorWorkflow в экспортированный воркфлоу n8n.
|
||||
// Использование: node n8n-fixerr.js <workflow.json> <errorWorkflowId>
|
||||
const fs = require('fs');
|
||||
let w = JSON.parse(fs.readFileSync(process.argv[2]));
|
||||
if (Array.isArray(w)) w = w[0];
|
||||
w.settings = w.settings || {};
|
||||
w.settings.errorWorkflow = process.argv[3];
|
||||
fs.writeFileSync(process.argv[2], JSON.stringify(w));
|
||||
console.log('OK ' + (w.name || w.id) + ' -> errorWorkflow=' + process.argv[3]);
|
||||
File diff suppressed because one or more lines are too long
@@ -0,0 +1,509 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
ПЦН-приёмник (ru3): Surgard MLR2-DG / Contact ID (:5111) + SIA DC-09 (:5112).
|
||||
Декодирует account + код события, публикует в mosquitto (raw MQTT 3.1.1, stdlib-only)
|
||||
топик guard/<account>/event JSON.
|
||||
|
||||
ПЕРВАЯ рабочая версия — минимальный, но живой приёмник.
|
||||
TODO: полная валидация SIA DC-09 CRC/шифрование; edge-cases Contact ID checksum.
|
||||
"""
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import socket
|
||||
import struct
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
import logging
|
||||
|
||||
# ---- Конфиг ----
|
||||
MQTT_HOST = "10.50.0.48" # mosquitto vm-mts1 по туннелю awg-mesh
|
||||
MQTT_PORT = 1883
|
||||
MQTT_KEEPALIVE = 60
|
||||
SURGARD_PORT = 5111
|
||||
SIA_PORT = 5112
|
||||
CLIENT_ID = "pcn-ru3"
|
||||
|
||||
# ---- Дисковый спул (буферизация при обрыве связи с брокером) ----
|
||||
SPOOL_DIR = "/var/spool/pcn"
|
||||
SPOOL_FILE = os.path.join(SPOOL_DIR, "queue.jsonl")
|
||||
SPOOL_MAX_BYTES = 100 * 1024 * 1024 # 100 МБ — лимит спула
|
||||
SPOOL_TRIM_BYTES = 80 * 1024 * 1024 # при ротации ужимаем до 80 МБ (дропаем старейшие)
|
||||
DRAIN_INTERVAL = 5 # секунд между попытками до-досыла
|
||||
DRAIN_BATCH = 200 # сколько сообщений максимум за один проход
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
format="%(asctime)s [%(levelname)s] %(message)s",
|
||||
)
|
||||
log = logging.getLogger("pcn")
|
||||
|
||||
# ---- Соответствие SIA <-> Surgard Contact ID (из src/data/security.ts) ----
|
||||
SIA_TO_CID = {
|
||||
"BA": "E130", "FA": "E110", "PA": "E120", "HA": "E121", "TA": "E137",
|
||||
"OP": "E401", "CL": "R401", "AT": "E301", "AR": "R301", "BT": "E302",
|
||||
"YC": "E354", "NL": "E381",
|
||||
}
|
||||
CID_TO_SIA = {v: k for k, v in SIA_TO_CID.items()}
|
||||
|
||||
|
||||
# ============================================================
|
||||
# Минимальный raw-MQTT publisher (3.1.1), без внешних зависимостей
|
||||
# ============================================================
|
||||
def _mqtt_remaining_length(n: int) -> bytes:
|
||||
out = bytearray()
|
||||
while True:
|
||||
b = n % 128
|
||||
n //= 128
|
||||
if n > 0:
|
||||
b |= 0x80
|
||||
out.append(b)
|
||||
if n == 0:
|
||||
break
|
||||
return bytes(out)
|
||||
|
||||
|
||||
def _mqtt_str(s: str) -> bytes:
|
||||
b = s.encode("utf-8")
|
||||
return struct.pack("!H", len(b)) + b
|
||||
|
||||
|
||||
class MqttPublisher:
|
||||
"""Синхронный публикатор в отдельном потоке-неблокирующий он не нужен:
|
||||
вызываем из executor, чтобы не блокировать asyncio-loop."""
|
||||
|
||||
def __init__(self, host, port, client_id, keepalive=60):
|
||||
self.host = host
|
||||
self.port = port
|
||||
self.client_id = client_id
|
||||
self.keepalive = keepalive
|
||||
self.sock = None
|
||||
# publish() и дренаж исполняются в разных потоках executor'а — сериализуем сокет.
|
||||
self._lock = threading.Lock()
|
||||
|
||||
def _connect(self):
|
||||
s = socket.create_connection((self.host, self.port), timeout=8)
|
||||
# CONNECT packet
|
||||
payload = _mqtt_str(self.client_id)
|
||||
var_header = _mqtt_str("MQTT") + bytes([0x04, 0x02]) + struct.pack("!H", self.keepalive)
|
||||
body = var_header + payload
|
||||
pkt = bytes([0x10]) + _mqtt_remaining_length(len(body)) + body
|
||||
s.sendall(pkt)
|
||||
# CONNACK (4 bytes)
|
||||
resp = s.recv(4)
|
||||
if len(resp) < 4 or resp[0] != 0x20 or resp[3] != 0x00:
|
||||
s.close()
|
||||
raise ConnectionError(f"MQTT CONNACK failed: {resp!r}")
|
||||
self.sock = s
|
||||
|
||||
def publish(self, topic, payload, qos=0):
|
||||
with self._lock:
|
||||
if self.sock is None:
|
||||
self._connect()
|
||||
try:
|
||||
self._publish(topic, payload)
|
||||
except (OSError, ConnectionError) as e:
|
||||
# reconnect once
|
||||
log.warning("MQTT publish err (%s), reconnecting", e)
|
||||
try:
|
||||
if self.sock:
|
||||
self.sock.close()
|
||||
except OSError:
|
||||
pass
|
||||
self.sock = None
|
||||
self._connect()
|
||||
self._publish(topic, payload)
|
||||
|
||||
def _publish(self, topic, payload):
|
||||
body = _mqtt_str(topic) + payload.encode("utf-8")
|
||||
pkt = bytes([0x30]) + _mqtt_remaining_length(len(body)) + body
|
||||
self.sock.sendall(pkt)
|
||||
|
||||
|
||||
_pub = MqttPublisher(MQTT_HOST, MQTT_PORT, CLIENT_ID, MQTT_KEEPALIVE)
|
||||
|
||||
|
||||
# ============================================================
|
||||
# Дисковый спул: append+fsync при обрыве, дренаж при восстановлении
|
||||
# ============================================================
|
||||
def _spool_ensure():
|
||||
try:
|
||||
os.makedirs(SPOOL_DIR, exist_ok=True)
|
||||
except OSError as e:
|
||||
log.error("spool mkdir err: %s", e)
|
||||
|
||||
|
||||
def _spool_append(topic, payload):
|
||||
"""Атомарно (append+fsync) дописать одно сообщение в очередь.
|
||||
Строка формата: {"topic":..,"payload":..}\n"""
|
||||
_spool_ensure()
|
||||
line = json.dumps({"topic": topic, "payload": payload}, ensure_ascii=False) + "\n"
|
||||
try:
|
||||
_spool_rotate_if_needed()
|
||||
with open(SPOOL_FILE, "a", encoding="utf-8") as f:
|
||||
f.write(line)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
log.warning("SPOOLED %s (queue on disk)", topic)
|
||||
except OSError as e:
|
||||
log.error("spool append err: %s (СООБЩЕНИЕ ПОТЕРЯНО)", e)
|
||||
|
||||
|
||||
def _spool_rotate_if_needed():
|
||||
"""Ограничение размера: если спул > SPOOL_MAX, дропаем старейшие строки
|
||||
до SPOOL_TRIM (сохраняем хвост = самые свежие события)."""
|
||||
try:
|
||||
if not os.path.exists(SPOOL_FILE):
|
||||
return
|
||||
size = os.path.getsize(SPOOL_FILE)
|
||||
if size <= SPOOL_MAX_BYTES:
|
||||
return
|
||||
log.warning("spool rotate: %d байт > лимит, ужимаю до %d", size, SPOOL_TRIM_BYTES)
|
||||
with open(SPOOL_FILE, "rb") as f:
|
||||
f.seek(max(0, size - SPOOL_TRIM_BYTES))
|
||||
f.readline() # выравниваемся на границу строки
|
||||
tail = f.read()
|
||||
tmp = SPOOL_FILE + ".tmp"
|
||||
with open(tmp, "wb") as f:
|
||||
f.write(tail)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
os.replace(tmp, SPOOL_FILE)
|
||||
except OSError as e:
|
||||
log.error("spool rotate err: %s", e)
|
||||
|
||||
|
||||
def _drain_once():
|
||||
"""Синхронно (в executor): попытаться до-дослать спул.
|
||||
Читает весь файл, публикует по одному; при первой ошибке — сохраняет
|
||||
непосланный остаток обратно (атомарно) и выходит. Возвращает кол-во посланных."""
|
||||
if not os.path.exists(SPOOL_FILE) or os.path.getsize(SPOOL_FILE) == 0:
|
||||
return 0
|
||||
sent = 0
|
||||
remainder = None
|
||||
try:
|
||||
with open(SPOOL_FILE, "r", encoding="utf-8") as f:
|
||||
lines = f.readlines()
|
||||
except OSError as e:
|
||||
log.error("drain read err: %s", e)
|
||||
return 0
|
||||
for i, line in enumerate(lines[:DRAIN_BATCH]):
|
||||
line = line.strip()
|
||||
if not line:
|
||||
sent += 1
|
||||
continue
|
||||
try:
|
||||
rec = json.loads(line)
|
||||
_pub.publish(rec["topic"], rec["payload"])
|
||||
sent += 1
|
||||
except (OSError, ConnectionError) as e:
|
||||
log.warning("drain: брокер снова недоступен (%s), стоп на %d", e, i)
|
||||
remainder = lines[i:]
|
||||
break
|
||||
except Exception as e:
|
||||
log.error("drain: битая строка спула, дроп: %s", e)
|
||||
sent += 1
|
||||
if remainder is None:
|
||||
remainder = lines[DRAIN_BATCH:]
|
||||
# Перезаписываем спул оставшимся хвостом атомарно
|
||||
try:
|
||||
if remainder:
|
||||
fd, tmp = tempfile.mkstemp(dir=SPOOL_DIR, prefix="q.", suffix=".tmp")
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as f:
|
||||
f.writelines(remainder)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
os.replace(tmp, SPOOL_FILE)
|
||||
else:
|
||||
os.remove(SPOOL_FILE)
|
||||
except OSError as e:
|
||||
log.error("drain rewrite err: %s", e)
|
||||
if sent:
|
||||
log.info("DRAINED %d из спула, осталось %d", sent, len(remainder))
|
||||
return sent
|
||||
|
||||
|
||||
async def drain_loop():
|
||||
"""Фоновая задача: раз в DRAIN_INTERVAL секунд пытается до-дослать спул."""
|
||||
loop = asyncio.get_event_loop()
|
||||
while True:
|
||||
await asyncio.sleep(DRAIN_INTERVAL)
|
||||
try:
|
||||
if os.path.exists(SPOOL_FILE) and os.path.getsize(SPOOL_FILE) > 0:
|
||||
await loop.run_in_executor(None, _drain_once)
|
||||
except Exception as e:
|
||||
log.error("drain_loop err: %s", e)
|
||||
|
||||
|
||||
async def publish_event(account, proto, cid, sia, zone, raw):
|
||||
doc = {
|
||||
"account": account,
|
||||
"proto": proto,
|
||||
"cid": cid,
|
||||
"sia": sia,
|
||||
"zone": zone,
|
||||
"raw": raw,
|
||||
"ts": int(time.time()),
|
||||
}
|
||||
topic = f"guard/{account}/event"
|
||||
payload = json.dumps(doc, ensure_ascii=False)
|
||||
loop = asyncio.get_event_loop()
|
||||
try:
|
||||
await loop.run_in_executor(None, _pub.publish, topic, payload)
|
||||
log.info("PUB %s %s", topic, payload)
|
||||
except Exception as e:
|
||||
# Брокер/туннель недоступны — НЕ теряем: спуливаем на диск.
|
||||
log.error("MQTT publish failed: %s -> спул", e)
|
||||
await loop.run_in_executor(None, _spool_append, topic, payload)
|
||||
|
||||
|
||||
# ============================================================
|
||||
# Разбор Contact ID кода -> (cid, sia)
|
||||
# ============================================================
|
||||
def cid_to_pair(qualifier, event3):
|
||||
"""qualifier '1'=new/E , '3'=restore/R ; event3 = 3 цифры (130,110,...)"""
|
||||
prefix = "R" if qualifier == "3" else "E"
|
||||
cid = f"{prefix}{event3}"
|
||||
sia = CID_TO_SIA.get(cid, "")
|
||||
return cid, sia
|
||||
|
||||
|
||||
# ============================================================
|
||||
# Surgard MLR2-DG / Contact ID приёмник (:5111)
|
||||
# ============================================================
|
||||
# Форматы, которые пытаемся распознать:
|
||||
# 1) Sur-Gard 5-2-1: "5BB 18 <account> 18 <CID> <qual> <zone>"
|
||||
# напр. "5011 18 1042 18 130 01 002"
|
||||
# 2) ADM-CID Contact ID: "<ACCT> 18 <MT> <Q> <XYZ> <GG> <CCC>"
|
||||
# где MT=18(DTMF)/98, Q=1/3, XYZ=код, GG=group/partition, CCC=zone/user
|
||||
SURGARD_ACK = b"\x06" # приёмник подтверждает событие символом ACK
|
||||
|
||||
|
||||
def parse_surgard(line: str):
|
||||
line = line.strip()
|
||||
if not line:
|
||||
return None
|
||||
toks = line.split()
|
||||
# Попытка: найти шаблон Contact ID где-то в строке.
|
||||
# Ищем последовательность: <acct> 18 <q><xyz> ... либо явные поля.
|
||||
# Нормализуем: убираем ведущий "5xx" код линии Sur-Gard если есть.
|
||||
# Универсальный разбор по токенам:
|
||||
# ищем токен "18" (Contact ID message type), после него qual+code.
|
||||
account = None
|
||||
cid = None
|
||||
sia = None
|
||||
zone = None
|
||||
|
||||
# Вариант A: пробельно-разделённые поля (тестовый формат)
|
||||
# "5011 18 1042 18 130 01 002"
|
||||
# idx: 0=lineprefix 1=18 2=account 3=18 4=code 5=group 6=zone
|
||||
if len(toks) >= 7 and toks[1] == "18" and toks[3] == "18":
|
||||
account = toks[2]
|
||||
code = toks[4] # 130
|
||||
# qualifier может идти как отдельная 1/3 или быть слит; тут по группе.
|
||||
# В классике group '01'/'02' — event; берём E по умолчанию.
|
||||
qual = "1"
|
||||
cid, sia = cid_to_pair(qual, code.zfill(3)[:3])
|
||||
zone = toks[6]
|
||||
return account, cid, sia, zone
|
||||
|
||||
# Вариант B: слитный Contact ID "ACCT 18 Q XYZ GG CCC"
|
||||
# напр. "1042 18 1 130 01 002"
|
||||
if len(toks) >= 6 and toks[1] in ("18", "98"):
|
||||
account = toks[0]
|
||||
qual = toks[2]
|
||||
code = toks[3]
|
||||
cid, sia = cid_to_pair(qual, code.zfill(3)[:3])
|
||||
zone = toks[-1]
|
||||
return account, cid, sia, zone
|
||||
|
||||
# Вариант C: полностью слитная ADM-CID строка "10421813001002"
|
||||
# ACCT(4) 18 Q(1) XYZ(3) GG(2) CCC(3)
|
||||
m = re.match(r"^(\d{3,4})18([13])(\d{3})(\d{2})(\d{3})$", line.replace(" ", ""))
|
||||
if m:
|
||||
account = m.group(1)
|
||||
qual = m.group(2)
|
||||
code = m.group(3)
|
||||
zone = m.group(5)
|
||||
cid, sia = cid_to_pair(qual, code)
|
||||
return account, cid, sia, zone
|
||||
|
||||
# Не распознали — вернём хотя бы raw с пустыми полями (account=UNKNOWN)
|
||||
return None
|
||||
|
||||
|
||||
async def handle_surgard(reader, writer):
|
||||
peer = writer.get_extra_info("peername")
|
||||
log.info("Surgard conn from %s", peer)
|
||||
try:
|
||||
while True:
|
||||
data = await asyncio.wait_for(reader.readline(), timeout=120)
|
||||
if not data:
|
||||
break
|
||||
raw = data.decode("latin-1", errors="replace").rstrip("\r\n")
|
||||
if not raw:
|
||||
continue
|
||||
log.info("Surgard raw: %r", raw)
|
||||
parsed = parse_surgard(raw)
|
||||
if parsed:
|
||||
account, cid, sia, zone = parsed
|
||||
await publish_event(account, "surgard", cid, sia, zone, raw)
|
||||
else:
|
||||
log.warning("Surgard: unparsed %r -> publishing raw under UNKNOWN", raw)
|
||||
await publish_event("UNKNOWN", "surgard", "", "", None, raw)
|
||||
# ACK приёмника
|
||||
writer.write(SURGARD_ACK)
|
||||
await writer.drain()
|
||||
except asyncio.TimeoutError:
|
||||
log.info("Surgard %s idle timeout", peer)
|
||||
except Exception as e:
|
||||
log.error("Surgard handler err: %s", e)
|
||||
finally:
|
||||
try:
|
||||
writer.close()
|
||||
except Exception:
|
||||
pass
|
||||
log.info("Surgard conn closed %s", peer)
|
||||
|
||||
|
||||
# ============================================================
|
||||
# SIA DC-09 приёмник (:5112)
|
||||
# ============================================================
|
||||
# Кадр DC-09 (упрощённо):
|
||||
# \x0A <CRC4hex> <LEN4hex> "<id>" <seq> R<rcvr> L<line> #<account> [<data>] <ts> \x0D
|
||||
# id: "SIA-DCS" | "ADM-CID" | "NULL" | ...
|
||||
# data для ADM-CID: |< contact-id ...> ; для SIA-DCS: |Nri1/BA01]
|
||||
# Ответ ACK: \x0A <CRC> <LEN> "ACK" <seq> R.. L.. #acct [] <ts> \x0D
|
||||
LF = 0x0A
|
||||
CR = 0x0D
|
||||
|
||||
|
||||
def _dc09_crc(msg: bytes) -> int:
|
||||
"""CRC-16 (polynomial 0x8005, reflected) как в DC-09 (для ACK)."""
|
||||
crc = 0
|
||||
for b in msg:
|
||||
crc ^= b
|
||||
for _ in range(8):
|
||||
if crc & 1:
|
||||
crc = (crc >> 1) ^ 0xA001
|
||||
else:
|
||||
crc >>= 1
|
||||
return crc & 0xFFFF
|
||||
|
||||
|
||||
def parse_sia(frame: str):
|
||||
"""frame — содержимое между LF и CR (без обрамляющих).
|
||||
Возвращает (account, cid, sia, zone, seq, id_token)."""
|
||||
account = None
|
||||
cid = None
|
||||
sia = None
|
||||
zone = None
|
||||
seq = "0000"
|
||||
id_token = ""
|
||||
|
||||
# id в кавычках
|
||||
m_id = re.search(r'"([A-Z\-]+)"', frame)
|
||||
if m_id:
|
||||
id_token = m_id.group(1)
|
||||
# seq — 4 цифры сразу после закрывающей кавычки id
|
||||
m_seq = re.search(r'"[A-Z\-]+"(\d{4})', frame)
|
||||
if m_seq:
|
||||
seq = m_seq.group(1)
|
||||
# account: #ACCT
|
||||
m_acct = re.search(r"#([0-9A-Fa-f]+)", frame)
|
||||
if m_acct:
|
||||
account = m_acct.group(1)
|
||||
|
||||
# payload между [ ... ] (у SIA-DCS) или после последнего |
|
||||
# SIA-DCS data: |Nri1/BA01] -> code=BA zone=01
|
||||
m_sia = re.search(r"/([A-Z]{2})(\d{1,3})", frame)
|
||||
if m_sia:
|
||||
sia = m_sia.group(1)
|
||||
zone = m_sia.group(2)
|
||||
cid = SIA_TO_CID.get(sia, "")
|
||||
else:
|
||||
# ADM-CID внутри DC-09: |1 XYZ GG CCC| или похожее
|
||||
m_cid = re.search(r"\|(\d)?\s?(\d{3})\s?(\d{2})?\s?(\d{2,3})?", frame)
|
||||
if m_cid:
|
||||
qual = m_cid.group(1) or "1"
|
||||
code = m_cid.group(2)
|
||||
zone = m_cid.group(4)
|
||||
prefix = "R" if qual == "3" else "E"
|
||||
cid = f"{prefix}{code}"
|
||||
sia = CID_TO_SIA.get(cid, "")
|
||||
|
||||
return account, cid, sia, zone, seq, id_token
|
||||
|
||||
|
||||
def build_sia_ack(seq: str, account: str) -> bytes:
|
||||
"""Строим минимальный DC-09 ACK-кадр."""
|
||||
ts = time.strftime("_%H:%M:%S,%m-%d-%Y", time.gmtime())
|
||||
acct = account or "0"
|
||||
body = f'"ACK"{seq}R0L0#{acct}[]{ts}'
|
||||
b = body.encode("ascii", errors="replace")
|
||||
crc = _dc09_crc(b)
|
||||
length = len(b)
|
||||
frame = f'\n{crc:04X}{length:04X}{body}\r'
|
||||
return frame.encode("ascii", errors="replace")
|
||||
|
||||
|
||||
async def handle_sia(reader, writer):
|
||||
peer = writer.get_extra_info("peername")
|
||||
log.info("SIA conn from %s", peer)
|
||||
try:
|
||||
while True:
|
||||
# читаем до CR (0x0D) — конец кадра DC-09
|
||||
data = await asyncio.wait_for(reader.readuntil(b"\r"), timeout=120)
|
||||
if not data:
|
||||
break
|
||||
raw = data.decode("latin-1", errors="replace")
|
||||
log.info("SIA raw: %r", raw)
|
||||
# вырезаем содержимое между LF и CR
|
||||
inner = raw
|
||||
if inner and inner[0] == "\n":
|
||||
inner = inner[1:]
|
||||
inner = inner.rstrip("\r")
|
||||
account, cid, sia, zone, seq, id_token = parse_sia(inner)
|
||||
if account is None:
|
||||
account = "UNKNOWN"
|
||||
await publish_event(account, "sia", cid, sia, zone, raw.strip())
|
||||
# ACK-кадр
|
||||
try:
|
||||
ack = build_sia_ack(seq, account if account != "UNKNOWN" else "0")
|
||||
writer.write(ack)
|
||||
await writer.drain()
|
||||
except Exception as e:
|
||||
log.warning("SIA ACK build/send err: %s", e)
|
||||
except asyncio.IncompleteReadError:
|
||||
log.info("SIA %s incomplete/closed", peer)
|
||||
except asyncio.TimeoutError:
|
||||
log.info("SIA %s idle timeout", peer)
|
||||
except Exception as e:
|
||||
log.error("SIA handler err: %s", e)
|
||||
finally:
|
||||
try:
|
||||
writer.close()
|
||||
except Exception:
|
||||
pass
|
||||
log.info("SIA conn closed %s", peer)
|
||||
|
||||
|
||||
# ============================================================
|
||||
async def main():
|
||||
srv_sur = await asyncio.start_server(handle_surgard, "0.0.0.0", SURGARD_PORT)
|
||||
srv_sia = await asyncio.start_server(handle_sia, "0.0.0.0", SIA_PORT)
|
||||
_spool_ensure()
|
||||
drain = asyncio.ensure_future(drain_loop())
|
||||
log.info("PCN receiver up: Surgard :%d, SIA DC-09 :%d -> mqtt %s:%d (spool %s)",
|
||||
SURGARD_PORT, SIA_PORT, MQTT_HOST, MQTT_PORT, SPOOL_FILE)
|
||||
async with srv_sur, srv_sia:
|
||||
await asyncio.gather(srv_sur.serve_forever(), srv_sia.serve_forever(), drain)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
try:
|
||||
asyncio.run(main())
|
||||
except KeyboardInterrupt:
|
||||
pass
|
||||
@@ -0,0 +1,16 @@
|
||||
[Unit]
|
||||
Description=PCN Receiver (Surgard :5111 + SIA DC-09 :5112) -> mosquitto vm-mts1
|
||||
After=network-online.target
|
||||
Wants=network-online.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
ExecStart=/usr/bin/python3 /opt/pcn-receiver/pcn-receiver.py
|
||||
Restart=always
|
||||
RestartSec=3
|
||||
User=root
|
||||
StandardOutput=journal
|
||||
StandardError=journal
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
File diff suppressed because one or more lines are too long
@@ -0,0 +1,296 @@
|
||||
"""Публикация телеметрии в MQTT с авто-обнаружением Home Assistant.
|
||||
|
||||
Каждый трекер (по IMEI) автоматически появляется в Home Assistant как
|
||||
устройство со своими сущностями: напряжение питания/АКБ, температуры,
|
||||
вскрытие шкафа, входы, уровень GSM, код события, позиция на карте.
|
||||
HA дальше используется как движок сцен/автоматизаций и агро-панель.
|
||||
|
||||
БУФЕРИЗАЦИЯ: при обрыве связи с брокером (туннель/mosquitto down) сообщения
|
||||
не теряются — пишутся в дисковый спул /opt/monitor-flex-ru2/spool/queue.jsonl
|
||||
(том, переживает пересборку), фоновый поток до-досылает при восстановлении.
|
||||
"""
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
|
||||
import paho.mqtt.client as mqtt
|
||||
|
||||
from . import config
|
||||
|
||||
log = logging.getLogger("flex.mqtt")
|
||||
|
||||
# ---- Дисковый спул ----
|
||||
SPOOL_DIR = os.getenv("SPOOL_DIR", "/app/spool")
|
||||
SPOOL_FILE = os.path.join(SPOOL_DIR, "queue.jsonl")
|
||||
SPOOL_MAX_BYTES = 100 * 1024 * 1024 # 100 МБ — лимит спула
|
||||
SPOOL_TRIM_BYTES = 80 * 1024 * 1024 # при ротации ужимаем до 80 МБ (дропаем старейшие)
|
||||
DRAIN_INTERVAL = 5 # секунд между попытками до-досыла
|
||||
DRAIN_BATCH = 500 # сколько сообщений максимум за один проход
|
||||
|
||||
# Числовые сенсоры: (ключ, имя, ед.изм., device_class, иконка)
|
||||
# Набор подобран под START S-2013: напряжение основного питания снимается с AIN3
|
||||
# (см. README, профиль S-2013), pwr_int = «системное питание» ≈ АКБ модуля.
|
||||
# Температуры приходят с RS-485/Bluetooth-датчиков (если настроены в конфигураторе).
|
||||
SENSORS = [
|
||||
("pwr_ext", "Напряжение питания (AIN3→Ug+)", "V", "voltage", None),
|
||||
("pwr_int", "Системное питание (АКБ)", "V", "voltage", None),
|
||||
("temp1", "Температура 1", "°C", "temperature", None),
|
||||
("temp2", "Температура 2", "°C", "temperature", None),
|
||||
("gsm", "Уровень GSM", None, None, "mdi:signal"),
|
||||
("event_code", "Код события", None, None, "mdi:alert-circle-outline"),
|
||||
]
|
||||
|
||||
# Бинарные датчики: (ключ, имя, device_class)
|
||||
# У S-2013 две дискретные линии (IN1, IN2); AIN3 занят под измерение напряжения.
|
||||
# Датчик двери/вскрытия шкафа подключается к IN1 или IN2.
|
||||
BINARY_SENSORS = [
|
||||
("in1", "Дверь шкафа (IN1)", "door"),
|
||||
("in2", "Охрана (IN2)", "safety"),
|
||||
("chasis_open", "Вскрытие корпуса", "tamper"),
|
||||
]
|
||||
|
||||
|
||||
class MqttPublisher:
|
||||
def __init__(self):
|
||||
self.client = mqtt.Client(client_id=config.MQTT_CLIENT_ID, protocol=mqtt.MQTTv311)
|
||||
if config.MQTT_USER:
|
||||
self.client.username_pw_set(config.MQTT_USER, config.MQTT_PASSWORD)
|
||||
self._announced = set()
|
||||
self._spool_lock = threading.Lock()
|
||||
self._connected = False
|
||||
self.client.on_connect = self._on_connect
|
||||
self.client.on_disconnect = self._on_disconnect
|
||||
self._drain_stop = threading.Event()
|
||||
self._drain_thread = None
|
||||
|
||||
# ---- paho callbacks: следим за реальным состоянием соединения ----
|
||||
def _on_connect(self, client, userdata, flags, rc):
|
||||
self._connected = (rc == 0)
|
||||
if rc == 0:
|
||||
log.info("MQTT connected (on_connect)")
|
||||
else:
|
||||
log.warning("MQTT connect rc=%s", rc)
|
||||
|
||||
def _on_disconnect(self, client, userdata, rc):
|
||||
self._connected = False
|
||||
if rc != 0:
|
||||
log.warning("MQTT disconnected rc=%s (буферизую на диск)", rc)
|
||||
|
||||
def connect(self):
|
||||
# connect_async + loop_start: не блокируемся, авто-reconnect средствами paho
|
||||
self.client.reconnect_delay_set(min_delay=1, max_delay=30)
|
||||
self.client.connect_async(config.MQTT_HOST, config.MQTT_PORT, keepalive=60)
|
||||
self.client.loop_start()
|
||||
self._spool_ensure()
|
||||
self._drain_thread = threading.Thread(target=self._drain_loop, daemon=True)
|
||||
self._drain_thread.start()
|
||||
log.info("MQTT loop запущен -> %s:%s (spool %s)",
|
||||
config.MQTT_HOST, config.MQTT_PORT, SPOOL_FILE)
|
||||
|
||||
def _state_topic(self, imei):
|
||||
return f"{config.MQTT_BASE_TOPIC}/{imei}/state"
|
||||
|
||||
def _device_block(self, imei):
|
||||
return {
|
||||
"identifiers": [f"navtelecom_{imei}"],
|
||||
"name": f"Щитовой шкаф {imei}",
|
||||
"manufacturer": "Navtelecom",
|
||||
"model": "FLEX tracker",
|
||||
}
|
||||
|
||||
def _announce(self, imei):
|
||||
"""Опубликовать конфиги обнаружения HA (один раз на устройство)."""
|
||||
if imei in self._announced:
|
||||
return
|
||||
state_topic = self._state_topic(imei)
|
||||
dev = self._device_block(imei)
|
||||
pfx = config.HA_DISCOVERY_PREFIX
|
||||
|
||||
for key, name, unit, dclass, icon in SENSORS:
|
||||
uid = f"navtelecom_{imei}_{key}"
|
||||
cfg = {
|
||||
"name": name,
|
||||
"unique_id": uid,
|
||||
"state_topic": state_topic,
|
||||
"value_template": f"{{{{ value_json.{key} }}}}",
|
||||
"device": dev,
|
||||
}
|
||||
if unit:
|
||||
cfg["unit_of_measurement"] = unit
|
||||
if dclass:
|
||||
cfg["device_class"] = dclass
|
||||
if dclass in ("voltage", "temperature"):
|
||||
cfg["state_class"] = "measurement"
|
||||
if icon:
|
||||
cfg["icon"] = icon
|
||||
self._pub(f"{pfx}/sensor/{uid}/config", cfg, retain=True)
|
||||
|
||||
for key, name, dclass in BINARY_SENSORS:
|
||||
uid = f"navtelecom_{imei}_{key}"
|
||||
cfg = {
|
||||
"name": name,
|
||||
"unique_id": uid,
|
||||
"state_topic": state_topic,
|
||||
"value_template": f"{{{{ value_json.{key} }}}}",
|
||||
"payload_on": 1,
|
||||
"payload_off": 0,
|
||||
"device_class": dclass,
|
||||
"device": dev,
|
||||
}
|
||||
self._pub(f"{pfx}/binary_sensor/{uid}/config", cfg, retain=True)
|
||||
|
||||
# Позиция на карте (device_tracker по GPS-атрибутам)
|
||||
uid = f"navtelecom_{imei}_tracker"
|
||||
self._pub(f"{pfx}/device_tracker/{uid}/config", {
|
||||
"name": f"Позиция {imei}",
|
||||
"unique_id": uid,
|
||||
"json_attributes_topic": state_topic,
|
||||
"source_type": "gps",
|
||||
"device": dev,
|
||||
}, retain=True)
|
||||
|
||||
self._announced.add(imei)
|
||||
log.info("HA discovery опубликован для %s", imei)
|
||||
|
||||
def publish(self, imei, record):
|
||||
"""Опубликовать одну запись телеметрии."""
|
||||
if not imei:
|
||||
return
|
||||
self._announce(imei)
|
||||
payload = dict(record)
|
||||
# device_tracker ждёт latitude/longitude
|
||||
if "lat" in record and "lon" in record:
|
||||
payload["latitude"] = record["lat"]
|
||||
payload["longitude"] = record["lon"]
|
||||
self._pub(self._state_topic(imei), payload, retain=True)
|
||||
|
||||
def _pub(self, topic, obj, retain=False):
|
||||
"""Публикация с буферизацией: если брокер недоступен — спуливаем на диск."""
|
||||
data = json.dumps(obj, ensure_ascii=False)
|
||||
self._pub_raw(topic, data, retain)
|
||||
|
||||
def _pub_raw(self, topic, data, retain=False):
|
||||
info = None
|
||||
if self._connected:
|
||||
try:
|
||||
info = self.client.publish(topic, data, retain=retain)
|
||||
except Exception as e:
|
||||
log.warning("MQTT publish exc: %s", e)
|
||||
info = None
|
||||
if info is None or info.rc != mqtt.MQTT_ERR_SUCCESS:
|
||||
# Брокер недоступен — не теряем, пишем в дисковый спул.
|
||||
self._spool_append(topic, data, retain)
|
||||
|
||||
# ============================================================
|
||||
# Дисковый спул: append+fsync при обрыве, дренаж при восстановлении
|
||||
# ============================================================
|
||||
def _spool_ensure(self):
|
||||
try:
|
||||
os.makedirs(SPOOL_DIR, exist_ok=True)
|
||||
except OSError as e:
|
||||
log.error("spool mkdir err: %s", e)
|
||||
|
||||
def _spool_append(self, topic, data, retain):
|
||||
line = json.dumps({"topic": topic, "payload": data, "retain": bool(retain)},
|
||||
ensure_ascii=False) + "\n"
|
||||
with self._spool_lock:
|
||||
try:
|
||||
self._spool_rotate_if_needed()
|
||||
with open(SPOOL_FILE, "a", encoding="utf-8") as f:
|
||||
f.write(line)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
log.warning("SPOOLED %s (брокер недоступен)", topic)
|
||||
except OSError as e:
|
||||
log.error("spool append err: %s (СООБЩЕНИЕ ПОТЕРЯНО)", e)
|
||||
|
||||
def _spool_rotate_if_needed(self):
|
||||
"""Ограничение размера: дропаем старейшие строки, сохраняя свежий хвост."""
|
||||
try:
|
||||
if not os.path.exists(SPOOL_FILE):
|
||||
return
|
||||
size = os.path.getsize(SPOOL_FILE)
|
||||
if size <= SPOOL_MAX_BYTES:
|
||||
return
|
||||
log.warning("spool rotate: %d байт > лимит, ужимаю до %d", size, SPOOL_TRIM_BYTES)
|
||||
with open(SPOOL_FILE, "rb") as f:
|
||||
f.seek(max(0, size - SPOOL_TRIM_BYTES))
|
||||
f.readline()
|
||||
tail = f.read()
|
||||
tmp = SPOOL_FILE + ".tmp"
|
||||
with open(tmp, "wb") as f:
|
||||
f.write(tail)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
os.replace(tmp, SPOOL_FILE)
|
||||
except OSError as e:
|
||||
log.error("spool rotate err: %s", e)
|
||||
|
||||
def _drain_once(self):
|
||||
with self._spool_lock:
|
||||
if not os.path.exists(SPOOL_FILE) or os.path.getsize(SPOOL_FILE) == 0:
|
||||
return 0
|
||||
try:
|
||||
with open(SPOOL_FILE, "r", encoding="utf-8") as f:
|
||||
lines = f.readlines()
|
||||
except OSError as e:
|
||||
log.error("drain read err: %s", e)
|
||||
return 0
|
||||
sent = 0
|
||||
remainder = None
|
||||
for i, line in enumerate(lines[:DRAIN_BATCH]):
|
||||
line = line.strip()
|
||||
if not line:
|
||||
sent += 1
|
||||
continue
|
||||
if not self._connected:
|
||||
remainder = lines[i:]
|
||||
break
|
||||
try:
|
||||
rec = json.loads(line)
|
||||
info = self.client.publish(rec["topic"], rec["payload"],
|
||||
retain=rec.get("retain", False))
|
||||
if info.rc != mqtt.MQTT_ERR_SUCCESS:
|
||||
remainder = lines[i:]
|
||||
break
|
||||
sent += 1
|
||||
except Exception as e:
|
||||
log.error("drain: битая строка спула, дроп: %s", e)
|
||||
sent += 1
|
||||
if remainder is None:
|
||||
remainder = lines[DRAIN_BATCH:]
|
||||
# Перезаписываем спул оставшимся хвостом атомарно
|
||||
with self._spool_lock:
|
||||
try:
|
||||
if remainder:
|
||||
fd, tmp = tempfile.mkstemp(dir=SPOOL_DIR, prefix="q.", suffix=".tmp")
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as f:
|
||||
f.writelines(remainder)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
os.replace(tmp, SPOOL_FILE)
|
||||
else:
|
||||
if os.path.exists(SPOOL_FILE):
|
||||
os.remove(SPOOL_FILE)
|
||||
except OSError as e:
|
||||
log.error("drain rewrite err: %s", e)
|
||||
if sent:
|
||||
log.info("DRAINED %d из спула, осталось %d", sent, len(remainder))
|
||||
return sent
|
||||
|
||||
def _drain_loop(self):
|
||||
while not self._drain_stop.is_set():
|
||||
self._drain_stop.wait(DRAIN_INTERVAL)
|
||||
try:
|
||||
if os.path.exists(SPOOL_FILE) and os.path.getsize(SPOOL_FILE) > 0:
|
||||
self._drain_once()
|
||||
except Exception as e:
|
||||
log.error("drain_loop err: %s", e)
|
||||
|
||||
def close(self):
|
||||
self._drain_stop.set()
|
||||
self.client.loop_stop()
|
||||
self.client.disconnect()
|
||||
@@ -0,0 +1 @@
|
||||
W1VuaXRdCkRlc2NyaXB0aW9uPVBDTiBSZWNlaXZlciAoU3VyZ2FyZCA6NTExMSArIFNJQSBEQy0wOSA6NTExMikgLT4gbW9zcXVpdHRvIHZtLW10czEKQWZ0ZXI9bmV0d29yay1vbmxpbmUudGFyZ2V0CldhbnRzPW5ldHdvcmstb25saW5lLnRhcmdldAoKW1NlcnZpY2VdClR5cGU9c2ltcGxlCkV4ZWNTdGFydD0vdXNyL2Jpbi9weXRob24zIC9vcHQvcGNuLXJlY2VpdmVyL3Bjbi1yZWNlaXZlci5weQpSZXN0YXJ0PWFsd2F5cwpSZXN0YXJ0U2VjPTMKVXNlcj1yb290ClN0YW5kYXJkT3V0cHV0PWpvdXJuYWwKU3RhbmRhcmRFcnJvcj1qb3VybmFsCgpbSW5zdGFsbF0KV2FudGVkQnk9bXVsdGktdXNlci50YXJnZXQK
|
||||
Reference in New Issue
Block a user