/
alex.piter
/
zigbridge
Обзор
Документация
Войти
/
alex.piter
/
zigbridge
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
data/web.py
641 строка
23 KB
Alexey Lapko
Добавил фильтрацию и сортировку ha_devices
28 май 2026, 13:45
28 май 2026, 13:45
fcc10bc
Код
Авторство
О чём код?
from flask import Flask, jsonify, request, render_template, send_from_directory import threading import db import os import json import logging import time import queue import serial import sys import serial.tools.list_ports import websocket from flask import Response, stream_with_context CONFIG_PATH="/data/options.json" logger = None level = logging.INFO level_mapping = {} if sys.version_info >= (3, 11): level_mapping = logging.getLevelNamesMapping() else: # Fallback for Python < 3.11 level_mapping = { logging.getLevelName(level): level for level in [ logging.CRITICAL, logging.ERROR, logging.WARNING, logging.INFO, logging.DEBUG, logging.NOTSET ] } try: with open(CONFIG_PATH, 'r') as file: config = json.load(file) if isinstance(config, dict) and 'log_level' in config: level = level_mapping.get(config.get('log_level')) except FileNotFoundError: print("The config file was not found.") except json.JSONDecodeError: print("The config file contains invalid JSON.") logging.basicConfig(level=level) logger = logging.getLogger(__name__) app = Flask(__name__, template_folder=os.path.join(os.path.dirname(__file__), '..', 'templates'), static_folder=os.path.join(os.path.dirname(__file__), '..', 'static')) # Initialize DB db.init() # WebSocket manager WS_URL = os.environ.get('WS_URL', 'ws://supervisor/core/websocket') class WebSocketManager: def __init__(self, url, token=None): self.url = url self.token = token self.ws = None self.thread = None self.connected = False self.out_queue = queue.Queue() self.pending = {} # id -> Queue self.lock = threading.Lock() def start(self): if self.thread and self.thread.is_alive(): return self.thread = threading.Thread(target=self._run, daemon=True) self.thread.start() def _run(self): backoff = 1 while True: try: self.ws = websocket.WebSocketApp(self.url, on_message=self._on_message, on_open=self._on_open, on_close=self._on_close) self.ws.run_forever( ping_interval=5 ) except Exception: logger.exception('WebSocket run_forever error') self.connected = False time.sleep(backoff) backoff = min(30, backoff * 2) def _on_open(self, ws): logger.info('WebSocket connected') self.connected = True # flush queued messages try: while not self.out_queue.empty(): msg = self.out_queue.get_nowait() try: ws.send(msg) except Exception: logger.exception('Failed sending queued message') self.out_queue.put(msg) break except Exception: logger.exception('Error flushing out_queue') # perform auth if token available try: token = os.environ.get('SUPERVISOR_TOKEN') if token: auth = {'type': 'auth', 'access_token': token} try: ws.send(json.dumps(auth)) except Exception: logger.exception('Failed to send auth payload') # trigger HA devices refresh async try: threading.Thread(target=refresh_ha_devices, daemon=True).start() except Exception: logger.exception('Failed to start HA refresh thread') except Exception: logger.exception('Auth handling in _on_open failed') def _on_close(self, ws, close_status_code, close_msg): logger.info('WebSocket closed: %s %s', close_status_code, close_msg) self.connected = False def _on_message(self, ws, message): logger.debug('Received msg from ws %s', message) try: payload = json.loads(message) except Exception: logger.exception('Failed to parse message: %s', message) type = payload.get('type') if type is not None and type == 'event': event = payload.get('event') if event is not None and event['event_type'] == 'state_changed': data = event['data'] new_state = data['new_state'] device_id = new_state['entity_id'] state = new_state['state'] device = db.get_device(device_id) if device is not None: state_int = 1 if state == 'on' else 0 device['status'] = state_int # update DB logger.info('Updating DB for cluster %s: ha_device=%s status=%s enabled=%s', device['cluster_id'], device['ha_device'], device['status'], device['enabled']) db.change_device(device) # push event push_sse_event({'type': 'devices_changed', 'cluster_id': device['cluster_id'], 'status': device['status']}) #send update to COM cluster_id = device['cluster_id'] payload = f'set {cluster_id} {state_int}' write_to_serial(payload) return # switch to on # INFO:__main__:Received msg from ws {"type":"event","event":{"event_type":"state_changed", # "data":{"entity_id":"switch.0x00124b00251ca3ba_l5","old_state":{"entity_id":"switch.0x00124b00251ca3ba_l5", # "state":"off","attributes":{"friendly_name":"0x00124b00251ca3ba L5"},"last_changed":"2026-05-27T11:16:16.108645+00:00", # "last_reported":"2026-05-27T11:16:16.108645+00:00","last_updated":"2026-05-27T11:16:16.108645+00:00", # "context":{"id":"01KSMJCH0BNZ58WHA4HWA23AAZ","parent_id":null,"user_id":"4813d0e8beeb4e33af12245d263ff5b6"}}, # "new_state":{"entity_id":"switch.0x00124b00251ca3ba_l5","state":"on","attributes":{"friendly_name":"0x00124b00251ca3ba L5"}, # "last_changed":"2026-05-27T11:17:09.234756+00:00","last_reported":"2026-05-27T11:17:09.234756+00:00", # "last_updated":"2026-05-27T11:17:09.234756+00:00","context":{"id":"01KSMJE4XNVD214F7R11X6EW6Y", # "parent_id":null,"user_id":"4813d0e8beeb4e33af12245d263ff5b6"}}},"origin":"LOCAL","time_fired":"2026-05-27T11:17:09.234756+00:00", # "context":{"id":"01KSMJE4XNVD214F7R11X6EW6Y","parent_id":null,"user_id":"4813d0e8beeb4e33af12245d263ff5b6"}},"id":1} mid = payload.get('id') if mid is None: return with self.lock: q = self.pending.get(mid) if q: try: q.put_nowait(payload) except Exception: logger.exception('Failed to deliver pending message') def send_raw(self, message): txt = message if isinstance(message, str) else json.dumps(message) if self.connected and self.ws: try: self.ws.send(txt) return True except Exception: logger.exception('send_raw failed, queueing') try: self.out_queue.put(txt) except Exception: logger.exception('Failed to queue outgoing message') return False def send_and_wait(self, payload, timeout=5): if not isinstance(payload, dict) or 'id' not in payload: self.send_raw(payload) return None mid = payload['id'] q = queue.Queue(maxsize=1) with self.lock: self.pending[mid] = q try: self.send_raw(payload) try: resp = q.get(timeout=timeout) return resp except queue.Empty: logger.warning('Timeout waiting for ws response id=%s', mid) return None finally: with self.lock: self.pending.pop(mid, None) ws_manager = WebSocketManager(WS_URL, token=os.environ.get('SUPERVISOR_TOKEN')) ws_manager.start() # Global state last_msg_id = 1 ser = None device_init = False reboot_state = False ha_devices_cache = [] # Simple server-sent events queue for notifying frontend of changes sse_queue = queue.Queue() def write_to_serial(payload): if ser: try: logger.debug('Try to write "%s" to serial', payload) payload += '\n' written = ser.write(payload.encode('utf-8')) try: ser.flush() except Exception: pass logger.info('Wrote %s bytes to serial: %s', written, payload.strip()) except serial.SerialTimeoutException: logger.exception("Write operation timed out") except Exception: logger.exception('Failed to write set to serial') else: logger.warning('Serial unavailable while attempting to send %s', payload) def push_sse_event(payload: dict): try: sse_queue.put_nowait(payload) except Exception: logger.exception('Failed to enqueue SSE event') @app.route('/events') def events(): def gen(): while True: try: ev = sse_queue.get() yield 'data: %s\n\n' % json.dumps(ev) except GeneratorExit: break except Exception: logger.exception('SSE generator error') time.sleep(0.1) return Response(stream_with_context(gen()), mimetype='text/event-stream') def refresh_ha_devices(): """Fetch HA states via websocket and populate a simple cache of entities with friendly names.""" global ha_devices_cache try: payload = {'id': get_request_id(), 'type': 'get_states'} resp = ws_manager.send_and_wait(payload, timeout=5) if not resp: logger.warning('No response when fetching HA states') # ha_devices_cache = [{'entity_id': "111111111", 'name': "lamp"}, # {'entity_id': "222222222", 'name': "lamp2"}] return # resp can be a list of entity dicts, or a dict with 'result', # or sometimes strings. Normalize accordingly. ha_devices_cache = [] candidates = [] if isinstance(resp, dict) and 'result' in resp: candidates = resp.get('result') or [] elif isinstance(resp, list): candidates = resp else: # Unexpected shape: try to tolerate a single entity string or dict candidates = [resp] for ent in candidates: try: if isinstance(ent, str): entity_id = ent name = ent elif isinstance(ent, dict): entity_id = ent.get('entity_id') or ent.get('entityId') or ent.get('id') attrs = ent.get('attributes', {}) if isinstance(ent.get('attributes', {}), dict) else {} name = attrs.get('friendly_name') or ent.get('name') or entity_id else: # fallback to string-converted representation entity_id = str(ent) name = entity_id if entity_id: domain = entity_id.split(".", 1)[0] allowed_domains = ["switch", "button", "light"] if domain in allowed_domains: ha_devices_cache.append({'entity_id': entity_id, 'name': name}) except Exception: logger.exception('Failed to parse entity from response element: %s', ent) ha_devices_cache.sort(key=lambda x: x['name']) except Exception: logger.exception('Failed to refresh HA devices') @app.route('/api/ha_devices', methods=['GET']) def api_ha_devices(): # Refresh cache in background if empty if not ha_devices_cache: try: threading.Thread(target=refresh_ha_devices, daemon=True).start() except Exception: logger.exception('Failed to refresh HA devices in bg') return jsonify(ha_devices_cache) def get_request_id(): global last_msg_id q = last_msg_id last_msg_id += 1 return q def connect_ws(): auth_data = { 'type': 'auth', 'access_token': os.environ.get('SUPERVISOR_TOKEN') } logger.info('auth payload: %s', auth_data) try: resp = ws_manager.send_and_wait(auth_data, timeout=5) logger.info('auth resp: %s', resp) except Exception: logger.exception('auth failed') def ping_ws(): while True: time.sleep(5) try: ws_manager.send_raw(json.dumps({'type': 'ping'})) except Exception: logger.exception('Failed to send ping') def get_ha_devices_states(): try: subscribe_payload = {'id': get_request_id(), 'type': 'subscribe_events', 'event_type': 'state_changed'} ws_manager.send_raw(json.dumps(subscribe_payload)) except Exception: logger.exception('WS error while subscribing events') return get_devices_payload = {'id': get_request_id(), 'type': 'get_states'} try: devices_response = ws_manager.send_and_wait(get_devices_payload, timeout=5) except Exception: logger.exception('WS error while getting devices') return if not devices_response: logger.warning('No devices response received') return try: devices = devices_response except Exception: logger.exception('Failed to parse devices response') return # Not wiring devices into UI here; frontend requests API def connect_uart(): global ser ports = serial.tools.list_ports.comports() for port, desc, hwid in sorted(ports): logger.info('%s %s [%s]', port, desc, hwid) ser_port = None for port, desc, hwid in sorted(ports): if 'VID:PID=303A:1001' in hwid: ser_port = port break logger.info('Device found at: %s', ser_port) if ser_port is None: logger.warning('No matching serial device found — running without serial connection') ser = None return try: ser = serial.Serial(ser_port, 115200, timeout=1, write_timeout = 1) logger.info('Device init complete') except Exception as e: logger.exception('Failed to open serial port %s: %s', ser_port, e) ser = None def serial_reconnector(): """Background thread: try to reconnect serial if `ser` is None.""" backoff = 1 while True: time.sleep(max(0.5, backoff)) global ser if ser: backoff = 1 continue try: connect_uart() if ser: logger.info('Serial reconnected') backoff = 1 try: logger.info('Sending init to ESP after reconnect') # setup() except Exception: logger.exception('Failed to send init after reconnect') else: backoff = min(30, backoff * 2) except Exception: logger.exception('Error in serial_reconnector') backoff = min(30, backoff * 2) def background_worker(): global device_init, ser # threading.Thread(target=reboot_device_loop, daemon=True).start() while True: if not ser: time.sleep(1) continue try: b = ser.readline() if len(b) == 0: continue s1 = b.decode('utf-8').rstrip() except Exception: logger.exception('Error reading from serial - closing and scheduling reconnect') try: try: ser.close() except Exception: pass ser = None except Exception: logger.exception('Failed to close serial after error') time.sleep(1) continue check = s1.split('|') if len(check) == 3: logger.info('Got state change: %s', check) try: data = json.loads(check[1].replace("'", '"')) handle_serial_state(data) except Exception: logger.exception('Failed to parse state payload') else: logger.debug('Received from esp: %s', s1) def handle_serial_state(data): # data contains {'cl': <cluster>, 'st': <state>} try: cl = data.get('cl') st = data.get('st') except Exception: return rows = db.get_devices() device_id = None for r in rows: if r['cluster_id'] == cl: device_id = r['ha_device'] break r['status'] = False if st == 0 else True logger.info('Updating DB for cluster %s: ha_device=%s status=%s enabled=%s', r['cluster_id'], r['ha_device'], r['status'], r['enabled']) db.change_device({'cluster_id': r['cluster_id'], 'ha_device': r['ha_device'], 'status': r['status'], 'enabled': r['enabled']}) try: # verify write current = [x for x in db.get_devices() if x['cluster_id'] == r['cluster_id']] logger.info('Post-update DB row for cluster %s: %s', r['cluster_id'], current) except Exception: logger.exception('Failed to read back device after update') try: push_sse_event({'type': 'devices_changed', 'cluster_id': r['cluster_id'], 'status': r['status']}) except Exception: logger.exception('Failed to push SSE after serial update') break if device_id: emit_device_state(device_id, 'turn_off' if st == 0 else 'turn_on') else: logger.warning('No device mapped for cluster %s', cl) def emit_device_state(device_id, state): # Infer domain from entity_id prefix (e.g. 'light.kitchen' -> 'light', 'input_boolean.test' -> 'input_boolean') domain = device_id.split('.', 1)[0] if isinstance(device_id, str) and '.' in device_id else 'light' payload = {'id': get_request_id(), 'type': 'call_service', 'domain': domain, 'service': state, 'target': {'entity_id': device_id}} try: logger.info('Sending HA call_service payload: %s', payload) resp = ws_manager.send_and_wait(payload, timeout=5) logger.info('Emitted state %s to %s, resp=%s', state, device_id, resp) except Exception: logger.exception('Failed to emit state') def start_background_workers(): try: threading.Thread(target=connect_ws, daemon=True).start() except Exception: logger.exception('Failed to start connect_ws') # try: # threading.Thread(target=ping_ws, daemon=True).start() # except Exception: # logger.exception('Failed to start ping_ws') try: threading.Thread(target=get_ha_devices_states, daemon=True).start() except Exception: logger.exception('Failed to start get_ha_devices_states') try: connect_uart() threading.Thread(target=background_worker, daemon=True).start() # start serial reconnector thread # threading.Thread(target=serial_reconnector, daemon=True).start() except Exception: logger.exception('Failed to start serial workers') def _before_first_request(): start_background_workers() # Try to register the startup handler with Flask if available; otherwise start immediately. bf = getattr(app, 'before_first_request', None) if callable(bf): try: bf(_before_first_request) except Exception: # If registration fails, start workers directly start_background_workers() else: start_background_workers() # Minimal routes for the frontend @app.route('/') def index(): try: return render_template('index.html') except Exception: return 'Index not available', 500 # Device API endpoints @app.route('/api/devices', methods=['GET']) def api_get_devices(): try: rows = db.get_devices() logger.debug('api_get_devices returning %s rows', len(rows)) return jsonify(rows) except Exception: logger.exception('Failed to list devices') return jsonify([]), 500 @app.route('/api/devices/<int:cluster_id>', methods=['PUT']) # вызывается когда жмут на чек бокс активация def api_add_device(cluster_id): try: data = request.get_json(force=True) # cluster_id cannot be changed data['cluster_id'] = cluster_id # validate presence of fields if 'ha_device' not in data: return jsonify({'error': 'nothing to add'}), 400 db.add_device(data) payload = f'add {cluster_id}' write_to_serial(payload) try: push_sse_event({'type': 'devices_changed'}) except Exception: logger.exception('Failed to push SSE after add') return jsonify({'ok': True}) except Exception: logger.exception('Failed to add device') return jsonify({'error': 'failed'}), 500 @app.route('/api/devices/<int:cluster_id>', methods=['DELETE']) def api_delete_device(cluster_id): try: db.delete_device(cluster_id) # try: # global reboot_state # reboot_state = True # logger.info('Requested ESP reboot after delete; reboot_state set') # if not ser: # threading.Thread(target=setup, daemon=True).start() # except Exception: # logger.exception('Failed to request reboot after delete') payload = f'del {cluster_id}' write_to_serial(payload) try: push_sse_event({'type': 'devices_changed'}) except Exception: logger.exception('Failed to push SSE after delete') return jsonify({'ok': True}) except Exception: logger.exception('Failed to delete device') return jsonify({'error': 'failed'}), 500 @app.route('/api/devices/<int:cluster_id>', methods=['POST']) # вызывается когда меняют привязанные к каналу ha_device def api_update_device(cluster_id): try: data = request.get_json(force=True) # cluster_id cannot be changed data['cluster_id'] = cluster_id # validate presence of fields if 'ha_device' not in data: return jsonify({'error': 'nothing to update'}), 400 db.change_device(data) try: push_sse_event({'type': 'devices_changed'}) except Exception: logger.exception('Failed to push SSE after add') return jsonify({'ok': True}), 201 except Exception: logger.exception('Failed to update device') return jsonify({'error': 'failed'}), 500 # Control endpoints @app.route('/api/apply', methods=['GET']) def api_apply(): write_to_serial('commit') if __name__ == '__main__': # Start Flask development server when run directly for testing logger.info('Starting Flask development server on 0.0.0.0:5000') logger.info('Startup info: WS_URL=%s SUPERVISOR_TOKEN=%s', WS_URL, 'present' if os.environ.get('SUPERVISOR_TOKEN') else 'absent') app.run(host='0.0.0.0', port=5000)