Files
frigate-camera-control-bridge/frigate_camera_control_bridge/app.py
T

1885 lines
77 KiB
Python

import json
import logging
import os
import queue
import re
import signal
import socket
import struct
import threading
import time
import urllib.error
import urllib.parse
import urllib.request
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
VERSION = "1.0"
DEFAULT_CONFIG = "/etc/frigate-camera-control-bridge.json"
LOG = logging.getLogger("frigate-camera-control-bridge")
STREAM_URL_RE = re.compile(r"(?:rtsp|rtsps|http|https)://[^\s'\"<>]+", re.IGNORECASE)
from .i18n import DEFAULT_LANGUAGE, available_languages, normalize_language, text_for
from .ui import read_static_file, render_page
def mqtt_encode_string(value):
# Minimal MQTT v3.1.1 encoder. Avoiding a client dependency keeps the image
# small and avoids pip installs on offline homelab hosts.
data = value.encode("utf-8")
return struct.pack("!H", len(data)) + data
def mqtt_encode_remaining_length(length):
out = bytearray()
while True:
digit = length % 128
length //= 128
if length:
digit |= 0x80
out.append(digit)
if not length:
return bytes(out)
def read_exact(sock, size):
buf = bytearray()
while len(buf) < size:
part = sock.recv(size - len(buf))
if not part:
raise EOFError("socket closed")
buf.extend(part)
return bytes(buf)
def read_mqtt_packet(sock):
first = read_exact(sock, 1)[0]
multiplier = 1
remaining_length = 0
while True:
digit = read_exact(sock, 1)[0]
remaining_length += (digit & 127) * multiplier
if not digit & 128:
break
multiplier *= 128
if multiplier > 128 * 128 * 128:
raise ValueError("malformed MQTT remaining length")
return first, read_exact(sock, remaining_length)
def send_mqtt_packet(sock, packet_type_with_flags, payload=b""):
sock.sendall(
bytes([packet_type_with_flags])
+ mqtt_encode_remaining_length(len(payload))
+ payload
)
def load_config(path):
with open(path, "r", encoding="utf-8") as handle:
return json.load(handle)
def slug(value):
return "".join(ch if ch.isalnum() else "_" for ch in value.lower()).strip("_")
def env_bool(name, default=True):
value = os.environ.get(name)
if value is None:
return default
return value.lower() in ("1", "true", "yes", "on")
def env_list(name):
value = os.environ.get(name, "")
return [item.strip() for item in value.split(",") if item.strip()]
def env_secret(name):
# Support both classic env vars and *_FILE secret mounts. Kubernetes users
# will usually use env from Secret, while Docker/Podman users may prefer a
# mounted secret file. The direct env var wins when both are present.
value = os.environ.get(name)
if value:
return value
path = os.environ.get(f"{name}_FILE")
if not path:
return None
try:
return Path(path).read_text(encoding="utf-8").strip()
except OSError as exc:
raise RuntimeError(f"failed to read {name}_FILE={path}: {exc}") from exc
def language_from_accept_header(value):
# HTTP Accept-Language can contain regional tags and q weights. The GUI
# ships language catalogs by base code, so pl-PL and pl both resolve to pl.
if not value:
return DEFAULT_LANGUAGE
supported = available_languages()
matches = []
for index, raw_item in enumerate(value.split(",")):
parts = [item.strip() for item in raw_item.split(";") if item.strip()]
if not parts:
continue
tag = parts[0].lower().replace("_", "-")
if tag == "*":
continue
quality = 1.0
for param in parts[1:]:
if not param.lower().startswith("q="):
continue
try:
quality = float(param.split("=", 1)[1])
except ValueError:
quality = 0.0
if quality <= 0:
continue
base = tag.split("-", 1)[0]
if base in supported:
matches.append((-quality, index, base))
if not matches:
return DEFAULT_LANGUAGE
matches.sort()
return matches[0][2]
def config_from_env():
# Env mode is the preferred sidecar mode. It needs only a Frigate API URL,
# MQTT connection parameters, and camera credentials. The old JSON config is
# still accepted for manual runs and migration debugging.
frigate_name = os.environ.get("FRIGATE_NAME", "frigate")
language = normalize_language(os.environ.get("UI_LANGUAGE", DEFAULT_LANGUAGE))
frigate_url = os.environ.get("FRIGATE_URL")
if not frigate_url:
frigate_host = os.environ.get("FRIGATE_HOST")
if not frigate_host:
raise RuntimeError("FRIGATE_URL or FRIGATE_HOST is required in env mode")
frigate_url = f"http://{frigate_host}:5000"
mqtt_host = os.environ.get("MQTT_HOST")
if not mqtt_host:
raise RuntimeError("MQTT_HOST is required in env mode")
mqtt_user = env_secret("MQTT_USER")
mqtt_password = env_secret("MQTT_PASSWORD")
if not mqtt_user:
raise RuntimeError("MQTT_USER is required in env mode")
if not mqtt_password:
raise RuntimeError("MQTT_PASSWORD is required in env mode")
camera_user = env_secret("DEFAULT_CAMERA_USER")
camera_password = env_secret("DEFAULT_CAMERA_PASSWORD")
if not camera_user:
raise RuntimeError("DEFAULT_CAMERA_USER is required in env mode")
if not camera_password:
raise RuntimeError("DEFAULT_CAMERA_PASSWORD is required in env mode")
mqtt = {
"host": mqtt_host,
"port": int(os.environ.get("MQTT_PORT", "1883")),
"base_topic": os.environ.get("BASE_TOPIC", f"frigate_camera_control/{slug(frigate_name)}"),
"discovery_prefix": os.environ.get("DISCOVERY_PREFIX", "homeassistant"),
"client_id": os.environ.get("MQTT_CLIENT_ID", f"frigate-camera-control-{slug(frigate_name)}"),
"username": mqtt_user,
"password": mqtt_password,
"keepalive": int(os.environ.get("MQTT_KEEPALIVE", "60")),
}
camera_auth = {
"username": camera_user,
"password": camera_password,
}
return {
"mqtt": mqtt,
"camera_auth": camera_auth,
"discovery": {
"enabled": True,
"state_file": os.environ.get(
"DISCOVERY_STATE_FILE",
"/var/lib/frigate-camera-control-bridge/discovery_state.json",
),
"include_hikvision_colorvu_white": env_bool("INCLUDE_HIKVISION_COLORVU_WHITE", True),
"include_dahua_coaxial_off_probe": env_bool("INCLUDE_DAHUA_COAXIAL_OFF_PROBE", True),
"wait_for_frigate": env_bool("WAIT_FOR_FRIGATE", True),
"retry_seconds": float(os.environ.get("FRIGATE_RETRY_SECONDS", "5")),
"startup_wait_seconds": float(os.environ.get("FRIGATE_STARTUP_WAIT_SECONDS", "0")),
"dahua_probe_channels": [int(item) for item in env_list("DAHUA_PROBE_CHANNELS") or ["1", "0", "2"]],
"dahua_siren_probe_channels": [
int(item) for item in env_list("DAHUA_SIREN_PROBE_CHANNELS") or ["2", "1", "0"]
],
"frigates": [
{
"name": frigate_name,
"host": os.environ.get("FRIGATE_HOST", urllib.parse.urlsplit(frigate_url).hostname or frigate_name),
"url": frigate_url,
}
],
},
"http": {
"enabled": env_bool("HTTP_ENABLED", True),
"bind": os.environ.get("HTTP_BIND", "0.0.0.0"),
"port": int(os.environ.get("HTTP_PORT", "5011")),
"token": env_secret("HTTP_TOKEN"),
},
"ui": {
"default_language": language,
},
"defaults": {
"http_timeout": float(os.environ.get("CAMERA_HTTP_TIMEOUT", "6")),
"colorvu_brightness": int(os.environ.get("COLORVU_BRIGHTNESS", "100")),
"snapshot_enabled": env_bool("SNAPSHOT_ENABLED", True),
"snapshot_height": int(os.environ.get("SNAPSHOT_HEIGHT", "180")),
"snapshot_cache_seconds": int(os.environ.get("SNAPSHOT_CACHE_SECONDS", "300")),
"snapshot_retry_seconds": int(os.environ.get("SNAPSHOT_RETRY_SECONDS", "600")),
"snapshot_refresh_seconds": int(os.environ.get("SNAPSHOT_REFRESH_SECONDS", "3600")),
"snapshot_live_on_gui": env_bool("SNAPSHOT_LIVE_ON_GUI", True),
"frigate_public_url": os.environ.get("FRIGATE_PUBLIC_URL"),
"frigate_camera_path_template": os.environ.get(
"FRIGATE_CAMERA_PATH_TEMPLATE",
"/#{camera_quoted}",
),
"frigate_refresh_seconds": int(os.environ.get("FRIGATE_REFRESH_SECONDS", "300")),
"auto_off_siren_seconds": int(os.environ.get("AUTO_OFF_SIREN_SECONDS", "180")),
"auto_off_red_blue_seconds": int(os.environ.get("AUTO_OFF_RED_BLUE_SECONDS", "300")),
"auto_off_white_light_seconds": int(os.environ.get("AUTO_OFF_WHITE_LIGHT_SECONDS", "600")),
"auto_off_default_seconds": int(os.environ.get("AUTO_OFF_DEFAULT_SECONDS", "0")),
"startup_all_off": env_bool("STARTUP_ALL_OFF", True),
},
"actions": [
{
"id": f"{slug(frigate_name)}_all_off",
"name": text_for(language, "action.all_off", target=frigate_name),
"name_key": "action.all_off",
"frigate_name": frigate_name,
"platform": "button",
"kind": "all_off",
"device": {
"identifiers": [f"frigate_camera_control_{slug(frigate_name)}"],
"name": f"Frigate camera control {frigate_name}",
"manufacturer": "Local",
"model": "Sidecar",
},
}
],
}
def load_config_or_env(path):
if path and Path(path).exists():
return load_config(path)
return config_from_env()
def response_status_ok(body):
text = body.decode("utf-8", errors="replace")
if "Error" in text and "ResponseStatus" not in text:
return False, text.strip()[:180]
if "<statusCode>1</statusCode>" in text or "<statusCode>0</statusCode>" in text:
return True, text.strip()[:180]
if "OK" in text and "Error" not in text:
return True, text.strip()[:180]
if not text.strip():
return True, ""
return True, text.strip()[:180]
class DigestHttpClient:
# Camera APIs in this installation use HTTP Digest. The client also allows
# Basic Auth because some camera firmwares accept either scheme.
def __init__(self, username, password, timeout=6):
self.username = username
self.password = password
self.timeout = timeout
def _opener(self, base_url):
mgr = urllib.request.HTTPPasswordMgrWithDefaultRealm()
mgr.add_password(None, base_url, self.username, self.password)
return urllib.request.build_opener(
urllib.request.HTTPDigestAuthHandler(mgr),
urllib.request.HTTPBasicAuthHandler(mgr),
)
def request(self, url, method="GET", body=None, content_type=None):
headers = {}
data = None
if body is not None:
data = body.encode("utf-8")
headers["Content-Type"] = content_type or "application/xml"
parsed = urllib.parse.urlsplit(url)
base_url = f"{parsed.scheme}://{parsed.netloc}/"
req = urllib.request.Request(url, data=data, headers=headers, method=method)
try:
with self._opener(base_url).open(req, timeout=self.timeout) as response:
payload = response.read()
ok, detail = response_status_ok(payload)
if not ok:
raise RuntimeError(detail)
return response.status, payload.decode("utf-8", errors="replace")
except urllib.error.HTTPError as exc:
payload = exc.read().decode("utf-8", errors="replace")
raise RuntimeError(f"HTTP {exc.code}: {payload[:180]}") from exc
except urllib.error.URLError as exc:
raise RuntimeError(str(exc.reason)) from exc
class Action:
# Action is the stable internal contract used by MQTT, HTTP GUI, CLI, and
# Home Assistant discovery. Camera-specific code is kept behind kind-specific
# runners so new vendors can be added without touching the MQTT layer.
def __init__(self, cfg, http_client):
self.cfg = cfg
self.id = cfg["id"]
self.name = cfg.get("name", self.id)
self.platform = cfg.get("platform", "switch")
self.kind = cfg["kind"]
self.http = http_client
self.state = None
self.auto_off_seconds = int(cfg.get("auto_off_seconds") or 0)
self.auto_off_at = None
self.snapshot_version = 0
def display_name(self, lang):
names = self.cfg.get("name_i18n") or {}
normalized = normalize_language(lang)
if names:
return names.get(normalized) or names.get(DEFAULT_LANGUAGE) or self.name
key = self.cfg.get("name_key")
if not key:
return self.name
prefix = self.cfg.get("camera_display_name") or self.cfg.get("camera_name") or self.cfg.get("frigate_name") or ""
return text_for(normalized, key, target=prefix).strip()
def command_list(self):
if self.kind == "all_off":
return ["press"]
if self.platform == "button":
return ["press"]
return ["on", "off"]
def safe_off_command(self):
if self.kind == "all_off":
return None
if "safe_off" in self.cfg:
return "safe_off"
if self.platform == "switch" and "off" in self.cfg:
return "off"
return None
def run(self, command, bridge=None):
command = command.lower()
if self.kind == "all_off":
if bridge is None:
raise RuntimeError("all_off requires bridge context")
return bridge.all_off()
if self.kind == "hikvision_colorvu":
self._run_hikvision(command)
elif self.kind == "dahua_coaxial_io":
self._run_dahua_coax(command)
else:
raise RuntimeError(f"unsupported action kind {self.kind}")
if self.platform == "switch":
if command == "on":
self.state = "ON"
elif command in ("off", "safe_off"):
self.state = "OFF"
return {"ok": True, "action": self.id, "command": command}
def _run_hikvision(self, command):
if command == "on":
self._hikvision_put_supplement(colorvu_xml(self.cfg.get("brightness", 100)))
self._hikvision_put_ircut(ircut_xml("night"))
return
if command in ("off", "safe_off"):
restore = self.cfg.get("restore", {})
self._hikvision_put_ircut(ircut_xml(restore.get("ircutFilterType", "auto")))
self._hikvision_put_supplement(restore_xml(self.cfg.get("brightness", 100), restore))
return
raise RuntimeError(f"unsupported command for {self.id}: {command}")
def _hikvision_put_supplement(self, xml):
host = self.cfg["host"]
self.http.request(
f"http://{host}/ISAPI/Image/channels/1/supplementLight",
method="PUT",
body=xml,
)
def _hikvision_put_ircut(self, xml):
host = self.cfg["host"]
self.http.request(
f"http://{host}/ISAPI/Image/channels/1/ircutFilter",
method="PUT",
body=xml,
)
def _run_dahua_coax(self, command):
params = self.cfg.get(command)
if params is None and command == "safe_off":
params = self.cfg.get("off") or self.cfg.get("press")
if params is None:
raise RuntimeError(f"unsupported command for {self.id}: {command}")
host = self.cfg["host"]
query = urllib.parse.urlencode(
{
"action": "control",
"channel": int(params["channel"]),
"info[0].Type": int(params["type"]),
"info[0].IO": int(params["io"]),
},
safe="[]",
)
self.http.request(f"http://{host}/cgi-bin/coaxialControlIO.cgi?{query}")
def colorvu_xml(brightness):
# Hikvision ColorVu white light is controlled through supplementLight plus
# ircutFilter. Turning it off restores the mode captured during discovery.
return f"""<?xml version="1.0" encoding="UTF-8"?>
<SupplementLight version="2.0" xmlns="http://www.hikvision.com/ver20/XMLSchema">
<supplementLightMode>colorVuWhiteLight</supplementLightMode>
<mixedLightBrightnessRegulatMode>manual</mixedLightBrightnessRegulatMode>
<whiteLightBrightness>{int(brightness)}</whiteLightBrightness>
<irLightBrightness>100</irLightBrightness>
<EventIntelligenceModeCfg>
<brightnessRegulatMode>manual</brightnessRegulatMode>
<whiteLightBrightness>{int(brightness)}</whiteLightBrightness>
<irLightBrightness>100</irLightBrightness>
</EventIntelligenceModeCfg>
</SupplementLight>"""
def restore_xml(brightness, restore):
associated = restore.get("associatedVMDHuman")
associated_xml = ""
if associated is not None:
associated_xml = f"\n<associatedVMDHuman>{str(bool(associated)).lower()}</associatedVMDHuman>"
return f"""<?xml version="1.0" encoding="UTF-8"?>
<SupplementLight version="2.0" xmlns="http://www.hikvision.com/ver20/XMLSchema">
<supplementLightMode>{restore.get("supplementLightMode", "eventIntelligence")}</supplementLightMode>
<mixedLightBrightnessRegulatMode>{restore.get("mixedLightBrightnessRegulatMode", "auto")}</mixedLightBrightnessRegulatMode>
<whiteLightBrightness>{int(restore.get("whiteLightBrightness", brightness))}</whiteLightBrightness>
<irLightBrightness>{int(restore.get("irLightBrightness", 100))}</irLightBrightness>
<EventIntelligenceModeCfg>
<brightnessRegulatMode>{restore.get("eventIntelligenceBrightnessRegulatMode", "auto")}</brightnessRegulatMode>
<whiteLightBrightness>{int(restore.get("whiteLightBrightness", brightness))}</whiteLightBrightness>
<irLightBrightness>{int(restore.get("irLightBrightness", 100))}</irLightBrightness>{associated_xml}
</EventIntelligenceModeCfg>
</SupplementLight>"""
def ircut_xml(mode):
return f"""<?xml version="1.0" encoding="UTF-8"?>
<IrcutFilter version="2.0" xmlns="http://www.hikvision.com/ver20/XMLSchema">
<IrcutFilterType>{mode}</IrcutFilterType>
<nightToDayFilterLevel>4</nightToDayFilterLevel>
<nightToDayFilterTime>5</nightToDayFilterTime>
</IrcutFilter>"""
class MqttBridgeClient:
# This is a deliberately small MQTT 3.1.1 client. It supports only what the
# bridge needs: connect, subscribe to command topics, publish retained HA
# discovery/state messages, and Last Will availability.
def __init__(self, cfg, messages):
self.cfg = cfg
self.messages = messages
self.sock = None
self.lock = threading.Lock()
self.packet_id = 1
@property
def base_topic(self):
return self.cfg.get("base_topic", "frigate_camera_control")
@property
def availability_topic(self):
return f"{self.base_topic}/status"
def next_packet_id(self):
self.packet_id += 1
if self.packet_id > 65535:
self.packet_id = 1
return self.packet_id
def connect(self, username=None, password=None):
host = self.cfg["host"]
port = int(self.cfg.get("port", 1883))
keepalive = int(self.cfg.get("keepalive", 60))
client_id = self.cfg.get("client_id", "frigate-camera-control-bridge")
sock = socket.create_connection((host, port), timeout=10)
sock.settimeout(1)
flags = 0x02
payload = mqtt_encode_string(client_id)
flags |= 0x04 | 0x20
payload += mqtt_encode_string(self.availability_topic)
payload += mqtt_encode_string("offline")
if username is not None:
flags |= 0x80
if password is not None:
flags |= 0x40
if username is not None:
payload += mqtt_encode_string(username)
if password is not None:
payload += mqtt_encode_string(password)
variable_header = mqtt_encode_string("MQTT") + bytes([4, flags]) + struct.pack("!H", keepalive)
send_mqtt_packet(sock, 0x10, variable_header + payload)
packet_type, body = read_mqtt_packet(sock)
if packet_type >> 4 != 2 or len(body) != 2 or body[1] != 0:
code = body[1] if len(body) >= 2 else "unknown"
raise RuntimeError(f"MQTT CONNACK failed: {code}")
discovery_prefix = self.cfg.get("discovery_prefix", "homeassistant")
topics = [f"{self.base_topic}/+/set", f"{discovery_prefix}/status"]
packet_id = self.next_packet_id()
subscribe_payload = struct.pack("!H", packet_id)
for topic in topics:
subscribe_payload += mqtt_encode_string(topic) + b"\x00"
send_mqtt_packet(sock, 0x82, subscribe_payload)
while True:
packet_type, body = read_mqtt_packet(sock)
if packet_type >> 4 != 9:
continue
if len(body) < 2 or struct.unpack("!H", body[:2])[0] != packet_id:
continue
if any(code == 0x80 for code in body[2:]):
raise RuntimeError(f"MQTT subscribe rejected for {topics}")
break
with self.lock:
self.sock = sock
return sock
def publish(self, topic, payload, retain=False):
if not isinstance(payload, bytes):
payload = str(payload).encode("utf-8")
packet = mqtt_encode_string(topic) + payload
with self.lock:
if self.sock is None:
raise RuntimeError("MQTT is not connected")
send_mqtt_packet(self.sock, 0x31 if retain else 0x30, packet)
def close(self):
with self.lock:
sock = self.sock
self.sock = None
if sock is not None:
try:
sock.close()
except OSError:
pass
def serve(self, credentials, on_connect, stop_event):
backoff = 1
credential_index = 0
keepalive = int(self.cfg.get("keepalive", 60))
while not stop_event.is_set():
credential = credentials[credential_index % len(credentials)]
credential_index += 1
sock = None
try:
sock = self.connect(credential.get("username"), credential.get("password"))
LOG.info("MQTT connected to %s:%s", self.cfg["host"], self.cfg.get("port", 1883))
on_connect()
backoff = 1
last_io = time.monotonic()
while not stop_event.is_set():
try:
first, body = read_mqtt_packet(sock)
except socket.timeout:
if time.monotonic() - last_io > max(10, keepalive / 2):
self.publish_ping()
last_io = time.monotonic()
continue
last_io = time.monotonic()
packet_type = first >> 4
flags = first & 0x0F
if packet_type == 3:
message = self.parse_publish(flags, body)
if message:
self.messages.put(message)
elif packet_type == 13:
continue
elif packet_type == 14:
raise EOFError("broker disconnected")
except Exception as exc:
if not stop_event.is_set():
LOG.warning("MQTT failed: %s", exc)
stop_event.wait(backoff)
backoff = min(backoff * 2, 60)
finally:
self.close()
if sock is not None:
try:
sock.close()
except OSError:
pass
def publish_ping(self):
with self.lock:
if self.sock is not None:
send_mqtt_packet(self.sock, 0xC0)
def disconnect_cleanly(self):
try:
self.publish(self.availability_topic, "offline", retain=True)
except Exception:
pass
with self.lock:
if self.sock is not None:
try:
send_mqtt_packet(self.sock, 0xE0)
except OSError:
pass
@staticmethod
def parse_publish(flags, body):
if len(body) < 2:
return None
topic_len = struct.unpack("!H", body[:2])[0]
if len(body) < 2 + topic_len:
return None
topic = body[2 : 2 + topic_len].decode("utf-8", errors="replace")
index = 2 + topic_len
qos = (flags >> 1) & 0x03
packet_id = None
if qos:
packet_id = struct.unpack("!H", body[index : index + 2])[0]
index += 2
payload = body[index:].decode("utf-8", errors="replace").strip()
return topic, payload, packet_id
class CameraControlBridge:
# CameraControlBridge wires together four boundaries:
# - Frigate API discovery,
# - camera vendor HTTP commands,
# - MQTT/Home Assistant discovery,
# - manual HTTP fallback GUI.
def __init__(self, config):
self.config = config
self.stop_event = threading.Event()
self.messages = queue.Queue()
self.action_lock = threading.RLock()
self.auto_off_timers = {}
self.credentials = self.load_mqtt_credentials()
camera_password = self.load_camera_password()
camera_user = config.get("camera_auth", {}).get("username", "admin")
self.default_language = normalize_language(
config.get("ui", {}).get("default_language", DEFAULT_LANGUAGE)
)
self.snapshot_cache = {}
self.http_client = DigestHttpClient(
camera_user,
camera_password,
timeout=float(config.get("defaults", {}).get("http_timeout", 6)),
)
self.inventory = []
self.actions = self.build_actions()
self.initialize_action_states()
self.refresh_snapshots()
self.mqtt = MqttBridgeClient(config["mqtt"], self.messages)
self.http_server = None
@property
def base_topic(self):
return self.config["mqtt"].get("base_topic", "frigate_camera_control")
def actions_snapshot(self):
with self.action_lock:
return list(self.actions.values())
def action_count(self):
with self.action_lock:
return len(self.actions)
def action_by_id(self, action_id):
with self.action_lock:
return self.actions.get(action_id)
def inherit_action_runtime_state(self, old_actions, new_actions):
for action_id, action in new_actions.items():
old = old_actions.get(action_id)
if old is None:
if action.platform == "switch":
action.state = "OFF"
continue
action.state = old.state
action.auto_off_at = old.auto_off_at
action.snapshot_version = old.snapshot_version
def initialize_action_states(self):
for action in self.actions.values():
if action.platform == "switch" and action.state is None:
action.state = "OFF"
def load_mqtt_credentials(self):
mqtt_cfg = self.config["mqtt"]
username = mqtt_cfg.get("username")
password = mqtt_cfg.get("password")
if not username:
raise RuntimeError("MQTT username is required")
if not password:
raise RuntimeError("MQTT password is required")
return [{"username": username, "password": password, "source": "environment"}]
def load_camera_password(self):
auth = self.config.get("camera_auth", {})
if auth.get("password"):
return auth["password"]
raise RuntimeError("default camera password is required")
def build_actions(self):
# Keep this bridge self-maintaining: every service start reads Frigate runtime
# configs, derives camera hosts, creates safe HA MQTT entities, and later
# removes stale discovery entities for cameras that disappeared.
self.inventory = self.discover_inventory()
present_hosts = {item["host"] for item in self.inventory if item.get("host")}
action_cfgs = []
for action_cfg in self.config.get("actions", []):
required_host = action_cfg.get("required_host")
if required_host and required_host not in present_hosts:
LOG.info("skip static action %s, host %s not in Frigate inventory", action_cfg["id"], required_host)
continue
action_cfgs.append(action_cfg)
action_cfgs.extend(self.discover_action_configs())
actions = {}
for action_cfg in action_cfgs:
cfg = dict(action_cfg)
auto_off_seconds = self.auto_off_seconds_for_config(cfg)
if auto_off_seconds > 0 and cfg.get("platform") == "switch":
cfg["auto_off_seconds"] = auto_off_seconds
if cfg["id"] in actions:
LOG.warning("duplicate action id skipped: %s", cfg["id"])
continue
actions[cfg["id"]] = Action(cfg, self.http_client)
return actions
def auto_off_seconds_for_config(self, cfg):
if "auto_off_seconds" in cfg:
return int(cfg.get("auto_off_seconds") or 0)
defaults = self.config.get("defaults", {})
name_key = cfg.get("name_key")
if name_key == "action.siren":
return int(defaults.get("auto_off_siren_seconds", 180))
if name_key == "action.red_blue":
return int(defaults.get("auto_off_red_blue_seconds", 300))
if name_key == "action.manual_white":
return int(defaults.get("auto_off_white_light_seconds", 600))
return int(defaults.get("auto_off_default_seconds", 0))
def discover_inventory(self):
# Frigate may reference cameras either through direct ffmpeg URLs or
# through local go2rtc gateway URLs. The selected source URLs are then
# used to discover the camera API host without hardcoded addresses.
discovery = self.config.get("discovery", {})
if not discovery.get("enabled", False):
return []
inventory = []
for frigate in discovery.get("frigates", []):
name = frigate.get("name") or frigate["host"]
base_url = frigate.get("url") or f"http://{frigate['host']}:5000"
data = self.load_frigate_config_with_retry(frigate, base_url)
if data is None:
continue
go2rtc = data.get("go2rtc") or {}
streams = go2rtc.get("streams") or go2rtc
cameras = data.get("cameras") or {}
for camera_name, camera_cfg in cameras.items():
urls = self.camera_candidate_urls(camera_name, camera_cfg, streams)
endpoints = self.camera_endpoints_from_urls(urls)
kinds = self.kinds_from_urls(urls)
camera_display_name = self.camera_display_name(camera_name, camera_cfg)
for host in endpoints:
inventory.append(
{
"frigate": name,
"frigate_url": base_url,
"camera": camera_name,
"camera_display_name": camera_display_name,
"host": host,
"kinds": kinds,
}
)
LOG.info("discovered %d camera host mapping(s) from Frigate", len(inventory))
return inventory
@staticmethod
def camera_display_name(camera_name, camera_cfg):
# Frigate camera keys remain the API identifier. friendly_name is only
# used for human-facing labels when present.
friendly_name = camera_cfg.get("friendly_name")
if friendly_name:
return str(friendly_name)
return camera_name
def load_frigate_config_with_retry(self, frigate, base_url):
discovery = self.config.get("discovery", {})
wait = bool(discovery.get("wait_for_frigate", True))
retry_seconds = float(discovery.get("retry_seconds", 5))
max_wait = float(discovery.get("startup_wait_seconds", 0))
started = time.monotonic()
source = frigate.get("config_file") or f"{base_url.rstrip('/')}/api/config"
while True:
try:
if frigate.get("config_file"):
return self.fetch_frigate_config_file(frigate["config_file"])
return self.fetch_json(f"{base_url.rstrip('/')}/api/config")
except Exception as exc:
if not wait:
LOG.warning("Frigate discovery failed for %s: %s", source, exc)
return None
elapsed = time.monotonic() - started
if max_wait and elapsed >= max_wait:
LOG.warning("Frigate discovery timed out for %s after %.0fs: %s", source, elapsed, exc)
return None
LOG.info("waiting for Frigate API %s: %s", source, exc)
time.sleep(retry_seconds)
@staticmethod
def fetch_json(url):
with urllib.request.urlopen(url, timeout=8) as response:
return json.loads(response.read().decode("utf-8", errors="replace"))
@staticmethod
def fetch_frigate_config_file(path):
try:
import yaml
except ImportError as exc:
raise RuntimeError("PyYAML is required for config_file discovery") from exc
with open(path, "r", encoding="utf-8") as handle:
return yaml.safe_load(handle) or {}
@staticmethod
def as_list(value):
if value is None:
return []
if isinstance(value, list):
return value
return [value]
@staticmethod
def camera_candidate_urls(camera_name, camera_cfg, streams):
# Discovery supports two Frigate layouts:
# 1. Direct camera URL in ffmpeg, e.g.
# rtsp://camera-a.example.test/...
# If such URL exists, it is the source of truth.
# 2. go2rtc gateway path in ffmpeg, e.g. rtsp://127.0.0.1:8554/stream-a.
# In this case the loopback URL is only a pointer; the real camera
# endpoint must be resolved from go2rtc.streams["stream-a"].
urls = CameraControlBridge.camera_direct_urls(camera_cfg)
if urls:
return urls
for stream_name in CameraControlBridge.camera_go2rtc_stream_names(camera_cfg):
resolved = CameraControlBridge.resolve_go2rtc_stream_urls(streams, stream_name)
if not resolved:
LOG.info("go2rtc stream %s/%s not found in Frigate config", camera_name, stream_name)
urls.extend(resolved)
return urls
@staticmethod
def camera_go2rtc_stream_names(camera_cfg):
# go2rtc-backed Frigate inputs usually look like:
# rtsp://127.0.0.1:8554/<stream_name>. Loopback here is not a camera API
# host; it is the go2rtc gateway. We keep only the stream name and
# resolve it through go2rtc.streams.
names = []
for item in camera_cfg.get("ffmpeg", {}).get("inputs", []):
path = str(item.get("path", ""))
stream_name = CameraControlBridge.go2rtc_stream_name_from_url(path)
if stream_name:
names.append(stream_name)
return list(dict.fromkeys(names))
@staticmethod
def camera_direct_urls(camera_cfg):
# Some installations skip go2rtc and put camera RTSP/HTTP URLs directly
# in ffmpeg inputs. Those URLs are inspected as-is. go2rtc loopback URLs
# are intentionally excluded here because they must be resolved through
# go2rtc.streams first.
urls = []
for item in camera_cfg.get("ffmpeg", {}).get("inputs", []):
path = str(item.get("path", ""))
if path and not CameraControlBridge.go2rtc_stream_name_from_url(path):
urls.extend(CameraControlBridge.stream_source_urls(path))
return urls
@staticmethod
def go2rtc_stream_name_from_url(url):
# go2rtc and Frigate often wrap real stream URLs, for example:
# ffmpeg:rtsp://127.0.0.1:8554/front-door#video=copy. Normalize that
# first so loopback gateway detection works for wrapped and raw values.
for candidate in CameraControlBridge.stream_source_urls(url) or [str(url)]:
try:
parsed = urllib.parse.urlsplit(candidate)
except ValueError:
continue
host = parsed.hostname or ""
if parsed.scheme != "rtsp" or parsed.port != 8554:
continue
if host not in ("127.0.0.1", "localhost", "::1"):
continue
return parsed.path.lstrip("/").split("/", 1)[0] or None
return None
@staticmethod
def stream_source_urls(value):
# go2rtc supports source wrappers such as "ffmpeg:rtsp://..." and may
# append options after a URL fragment. Extract URL tokens instead of
# treating the wrapper prefix as the network scheme.
text = str(value or "")
urls = STREAM_URL_RE.findall(text)
if urls:
return urls
if "://" in text:
return [text]
return []
@staticmethod
def resolve_go2rtc_stream_urls(streams, stream_name):
# go2rtc stream values may be a string, a list of strings, or a small
# mapping depending on how the Frigate config was authored. Flatten only
# the values that can contain source URLs; non-URL helper values are
# harmless because host/vendor extraction is conservative.
value = streams.get(stream_name) if isinstance(streams, dict) else None
return CameraControlBridge.flatten_go2rtc_value(value)
@staticmethod
def flatten_go2rtc_value(value):
if value is None:
return []
if isinstance(value, str):
urls = CameraControlBridge.stream_source_urls(value)
return urls or [value]
if isinstance(value, list):
out = []
for item in value:
out.extend(CameraControlBridge.flatten_go2rtc_value(item))
return out
if isinstance(value, dict):
out = []
for key in ("streams", "url", "urls", "source", "sources"):
if key in value:
out.extend(CameraControlBridge.flatten_go2rtc_value(value[key]))
return out
return [str(value)]
@staticmethod
def hosts_from_urls(urls):
return CameraControlBridge.camera_endpoints_from_urls(urls)
@staticmethod
def camera_endpoints_from_urls(urls):
# Extract the camera API endpoint from URLs without hardcoding camera IPs.
# IP literals and DNS names are both valid. Loopback addresses should
# already have been resolved through go2rtc.streams, so skip them if a
# config still leaks one through.
endpoints = []
for url in urls:
try:
parsed = urllib.parse.urlsplit(str(url))
except ValueError:
continue
host = parsed.hostname
if not host:
continue
normalized = host.strip("[]").lower()
if (
normalized in ("localhost", "::1", "0.0.0.0", "255.255.255.255")
or normalized.startswith("127.")
):
continue
endpoints.append(normalized)
return list(dict.fromkeys(endpoints))
@staticmethod
def kinds_from_urls(urls):
# Vendor detection is based on URL shape. This decides which safe
# capability probes may be attempted for each discovered camera.
kinds = set()
for url in urls:
if "Streaming/Channels" in url:
kinds.add("hikvision")
elif "cam/realmonitor" in url:
kinds.add("dahua")
elif "/stream1" in url or "/stream2" in url:
kinds.add("tapo")
elif "camera/stream" in url:
kinds.add("tablet")
else:
kinds.add("other")
return sorted(kinds)
def discover_action_configs(self):
discovery = self.config.get("discovery", {})
if not discovery.get("enabled", False):
return []
actions = []
for item in self.primary_inventory_items():
host = item["host"]
if discovery.get("include_hikvision_colorvu_white", True):
action = self.discover_hikvision_colorvu_action(item)
if action:
actions.append(action)
if discovery.get("include_dahua_coaxial_off_probe", True):
actions.extend(self.discover_dahua_coaxial_actions(item))
LOG.info("auto-created %d action(s) from camera capabilities", len(actions))
return actions
def primary_inventory_items(self):
grouped = {}
for item in self.inventory:
host = item["host"]
current = grouped.get(host)
if current is None:
grouped[host] = item
continue
if "ptz" in item["camera"].lower() and "ptz" not in current["camera"].lower():
grouped[host] = item
return list(grouped.values())
def discover_hikvision_colorvu_action(self, item):
if "hikvision" not in item.get("kinds", []):
return None
host = item["host"]
try:
_, device_info = self.http_client.request(f"http://{host}/ISAPI/System/deviceInfo")
model = self.xml_text(device_info, "model") or "Hikvision"
_, caps = self.http_client.request(f"http://{host}/ISAPI/Image/channels/1/supplementLight/capabilities")
modes = self.xml_attr(caps, "supplementLightMode", "opt") or ""
if "colorVuWhiteLight" not in modes.split(","):
return None
_, current = self.http_client.request(f"http://{host}/ISAPI/Image/channels/1/supplementLight")
except Exception as exc:
LOG.info("skip Hikvision discovery %s/%s: %s", item["camera"], host, exc)
return None
restore = {
"supplementLightMode": self.xml_text(current, "supplementLightMode") or "eventIntelligence",
"mixedLightBrightnessRegulatMode": self.xml_text(current, "mixedLightBrightnessRegulatMode") or "auto",
"eventIntelligenceBrightnessRegulatMode": self.xml_text(current, "brightnessRegulatMode") or "auto",
"whiteLightBrightness": int(self.xml_text(current, "whiteLightBrightness") or 100),
"irLightBrightness": int(self.xml_text(current, "irLightBrightness") or 100),
"ircutFilterType": "auto",
}
associated = self.xml_text(current, "associatedVMDHuman")
if associated is not None:
restore["associatedVMDHuman"] = associated.lower() == "true"
action_id = f"{slug(item['frigate'])}_{slug(item['camera'])}_colorvu_white"
lang = self.default_language
return {
"id": action_id,
"name": text_for(lang, "action.manual_white", target=item["camera"]),
"name_key": "action.manual_white",
"camera_name": item["camera"],
"camera_display_name": item.get("camera_display_name", item["camera"]),
"platform": "switch",
"kind": "hikvision_colorvu",
"host": host,
"brightness": int(self.config.get("defaults", {}).get("colorvu_brightness", 100)),
"restore": restore,
"snapshot": self.snapshot_config(item),
"device": {
"identifiers": [f"camera_{host.replace('.', '_')}"],
"name": item.get("camera_display_name", item["camera"]),
"manufacturer": "Hikvision",
"model": model,
"configuration_url": item["frigate_url"],
},
}
def discover_dahua_coaxial_actions(self, item):
# Dahua red-blue lights and sirens are discovered with OFF-only probes.
# A successful OFF endpoint means the command exists without starting
# the alarm/light during bridge startup.
if "dahua" not in item.get("kinds", []):
return []
host = item["host"]
try:
_, info = self.http_client.request(f"http://{host}/cgi-bin/magicBox.cgi?action=getSystemInfo")
except Exception as exc:
LOG.info("skip Dahua discovery %s/%s: %s", item["camera"], host, exc)
return []
model = self.dahua_value(info, "deviceType") or "Dahua"
actions = []
red_blue_channel = self.dahua_off_probe(host, action_type=1)
lang = self.default_language
if red_blue_channel is not None:
actions.append(
{
"id": f"{slug(item['frigate'])}_{slug(item['camera'])}_red_blue",
"name": text_for(lang, "action.red_blue", target=item["camera"]),
"name_key": "action.red_blue",
"camera_name": item["camera"],
"camera_display_name": item.get("camera_display_name", item["camera"]),
"platform": "switch",
"kind": "dahua_coaxial_io",
"host": host,
"on": {"channel": red_blue_channel, "type": 1, "io": 1},
"off": {"channel": red_blue_channel, "type": 1, "io": 2},
"snapshot": self.snapshot_config(item),
"device": self.discovered_device(item, "Dahua", model),
}
)
siren_channel = self.dahua_off_probe(host, action_type=2)
if siren_channel is not None:
actions.append(
{
"id": f"{slug(item['frigate'])}_{slug(item['camera'])}_siren",
"name": text_for(lang, "action.siren", target=item["camera"]),
"name_key": "action.siren",
"camera_name": item["camera"],
"camera_display_name": item.get("camera_display_name", item["camera"]),
"platform": "switch",
"kind": "dahua_coaxial_io",
"host": host,
"on": {"channel": siren_channel, "type": 2, "io": 1},
"off": {"channel": siren_channel, "type": 2, "io": 2},
"safe_off": {"channel": siren_channel, "type": 2, "io": 2},
"requires_ack": True,
"snapshot": self.snapshot_config(item),
"device": self.discovered_device(item, "Dahua", model),
}
)
return actions
def dahua_off_probe(self, host, action_type):
channels = self.config.get("discovery", {}).get("dahua_probe_channels", [1, 2, 0])
if action_type == 2:
channels = self.config.get("discovery", {}).get("dahua_siren_probe_channels", [2, 1, 0])
for channel in channels:
query = urllib.parse.urlencode(
{
"action": "control",
"channel": int(channel),
"info[0].Type": int(action_type),
"info[0].IO": 2,
},
safe="[]",
)
try:
self.http_client.request(f"http://{host}/cgi-bin/coaxialControlIO.cgi?{query}")
LOG.info("Dahua OFF probe OK host=%s type=%s channel=%s", host, action_type, channel)
return int(channel)
except Exception:
continue
return None
@staticmethod
def dahua_value(text, key):
match = re.search(rf"^{re.escape(key)}=(.+)$", text, re.MULTILINE)
return match.group(1).strip().strip("\r") if match else None
@staticmethod
def discovered_device(item, manufacturer, model):
return {
"identifiers": [f"camera_{item['host'].replace('.', '_')}"],
"name": item.get("camera_display_name", item["camera"]),
"manufacturer": manufacturer,
"model": model,
"configuration_url": item["frigate_url"],
}
def snapshot_config(self, item):
config = {
"frigate_url": item["frigate_url"],
"camera": item["camera"],
"camera_display_name": item.get("camera_display_name", item["camera"]),
}
defaults = self.config.get("defaults", {})
if defaults.get("frigate_public_url"):
config["public_url"] = defaults["frigate_public_url"]
if defaults.get("frigate_camera_path_template"):
config["camera_path_template"] = defaults["frigate_camera_path_template"]
return config
@staticmethod
def xml_text(xml, tag):
match = re.search(rf"<(?:[^:>]+:)?{re.escape(tag)}[^>]*>([^<]+)</(?:[^:>]+:)?{re.escape(tag)}>", xml)
return match.group(1).strip() if match else None
@staticmethod
def xml_attr(xml, tag, attr):
match = re.search(rf"<(?:[^:>]+:)?{re.escape(tag)}[^>]*\s{re.escape(attr)}=\"([^\"]*)\"", xml)
return match.group(1).strip() if match else None
def list_actions(self):
for action in self.actions_snapshot():
commands = ",".join(action.command_list())
print(f"{action.id}\t{action.platform}\t{commands}\t{action.display_name(self.default_language)}")
def check(self):
ok = True
for credential in self.credentials:
try:
test = MqttBridgeClient(self.config["mqtt"], queue.Queue())
sock = test.connect(credential.get("username"), credential.get("password"))
sock.close()
print(f"MQTT OK: {self.config['mqtt']['host']} as {credential['source']}")
break
except Exception as exc:
ok = False
print(f"MQTT ERROR using {credential['source']}: {exc}")
for action in self.actions_snapshot():
if action.kind == "hikvision_colorvu":
print(f"CHECK: {action.id} Hikvision host={action.cfg['host']} configured")
elif action.kind == "dahua_coaxial_io":
print(f"CHECK: {action.id} Dahua host={action.cfg['host']} configured")
elif action.kind == "all_off":
print(f"CHECK: {action.id} all-off button configured")
return 0 if ok else 1
def refresh_snapshots(self):
# Snapshot refresh is best-effort. It improves the manual GUI but should
# never prevent MQTT controls from starting.
if not self.config.get("defaults", {}).get("snapshot_enabled", True):
return True
seen = set()
ok = True
for action in self.actions_snapshot():
url = self.snapshot_url(action)
if not url or url in seen:
continue
seen.add(url)
try:
self.fetch_snapshot(url, force=True)
except Exception as exc:
ok = False
LOG.info("snapshot unavailable for %s: %s", action.id, exc)
return ok
@staticmethod
def should_refresh_snapshot_after_command(action, command):
if command.lower() not in ("on", "off", "safe_off"):
return False
return action.cfg.get("name_key") in ("action.manual_white", "action.red_blue")
def invalidate_action_snapshot(self, action):
url = self.snapshot_url(action)
if not url:
return
with self.action_lock:
self.snapshot_cache.pop(url, None)
version = int(time.time() * 1000)
for item in self.actions.values():
if self.snapshot_url(item) == url:
item.snapshot_version = version
def snapshot_url(self, action):
snapshot = action.cfg.get("snapshot")
if not snapshot:
return None
camera = urllib.parse.quote(snapshot["camera"], safe="")
height = int(self.config.get("defaults", {}).get("snapshot_height", 180))
return f"{snapshot['frigate_url'].rstrip('/')}/api/{camera}/latest.jpg?h={height}"
def fetch_snapshot(self, url, force=False):
now = time.monotonic()
max_age = int(self.config.get("defaults", {}).get("snapshot_cache_seconds", 300))
cached = self.snapshot_cache.get(url)
if cached and not force and now - cached["time"] < max_age:
return cached
with urllib.request.urlopen(url, timeout=8) as response:
data = response.read()
content_type = response.headers.get("Content-Type", "image/jpeg").split(";", 1)[0]
if not data:
raise RuntimeError("empty snapshot")
cached = {"data": data, "content_type": content_type, "time": now}
self.snapshot_cache[url] = cached
return cached
def snapshot_for_action(self, action_id, live=False):
action = self.action_by_id(action_id)
if not action:
return None
url = self.snapshot_url(action)
if not url:
return None
try:
return self.fetch_snapshot(url, force=live)
except Exception:
cached = self.snapshot_cache.get(url)
if cached:
return cached
raise
def publish_discovery(self):
# Home Assistant MQTT discovery topics are retained. The bridge remembers
# what it published previously and clears stale entities when Frigate no
# longer exposes a camera/action.
discovery_prefix = self.config["mqtt"].get("discovery_prefix", "homeassistant")
availability = f"{self.base_topic}/status"
self.mqtt.publish(availability, "online", retain=True)
previous = self.load_discovery_state()
current = {}
for action in self.actions_snapshot():
payload = self.discovery_payload(action, availability)
topic = f"{discovery_prefix}/{action.platform}/frigate_camera_control/{action.id}/config"
self.mqtt.publish(topic, json.dumps(payload, separators=(",", ":")), retain=True)
self.publish_state(action, retain=True)
current[action.id] = {"platform": action.platform, "topic": topic}
for action_id, old in previous.items():
if action_id not in current:
topic = old.get("topic") or f"{discovery_prefix}/{old.get('platform', 'switch')}/frigate_camera_control/{action_id}/config"
self.mqtt.publish(topic, "", retain=True)
LOG.info("removed stale MQTT discovery entity %s", action_id)
self.save_discovery_state(current)
LOG.info("published MQTT discovery for %d actions", len(current))
def publish_state(self, action, retain=True):
if action.platform == "switch" and action.state in ("ON", "OFF"):
self.mqtt.publish(f"{self.base_topic}/{action.id}/state", action.state, retain=retain)
def discovery_state_file(self):
return Path(
self.config.get("discovery", {}).get(
"state_file",
"/var/lib/frigate-camera-control-bridge/discovery_state.json",
)
)
def load_discovery_state(self):
path = self.discovery_state_file()
try:
return json.loads(path.read_text(encoding="utf-8"))
except FileNotFoundError:
return {}
except Exception as exc:
LOG.warning("failed to read discovery state %s: %s", path, exc)
return {}
def save_discovery_state(self, state):
path = self.discovery_state_file()
try:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(state, indent=2), encoding="utf-8")
except Exception as exc:
LOG.warning("failed to write discovery state %s: %s", path, exc)
def discovery_payload(self, action, availability_topic):
base = f"{self.base_topic}/{action.id}"
device_cfg = action.cfg.get("device") or {}
device = {
"identifiers": device_cfg.get("identifiers", ["frigate_camera_control"]),
"name": device_cfg.get("name", "Frigate Camera Control"),
}
for key in ("manufacturer", "model", "sw_version", "configuration_url"):
if device_cfg.get(key):
device[key] = device_cfg[key]
payload = {
"name": action.display_name(self.default_language),
"unique_id": f"frigate_camera_control_{action.id}",
"availability_topic": availability_topic,
"payload_available": "online",
"payload_not_available": "offline",
"device": device,
"origin": {
"name": "frigate-camera-control-bridge",
"sw_version": VERSION,
},
"qos": 0,
}
if action.platform == "switch":
payload.update(
{
"command_topic": f"{base}/set",
"state_topic": f"{base}/state",
"payload_on": "ON",
"payload_off": "OFF",
"state_on": "ON",
"state_off": "OFF",
"optimistic": False,
"retain": False,
}
)
elif action.platform == "button":
payload.update(
{
"command_topic": f"{base}/set",
"payload_press": "PRESS",
"retain": False,
}
)
return payload
def execute_action(self, action_id, command, publish=True):
action = self.action_by_id(action_id)
if action is None:
raise RuntimeError(f"unknown action {action_id}")
if command == "on" and self.auto_off_running(action):
raise RuntimeError("action is already active until automatic turn-off")
result = action.run(command, bridge=self)
if action.platform == "switch":
if command == "on":
self.schedule_auto_off(action)
elif command in ("off", "safe_off"):
self.cancel_auto_off(action.id)
if self.should_refresh_snapshot_after_command(action, command):
self.invalidate_action_snapshot(action)
if publish:
try:
self.publish_result(action, command, result)
except Exception as exc:
LOG.warning("action %s %s worked, but MQTT publish failed: %s", action_id, command, exc)
return result
@staticmethod
def auto_off_running(action):
return bool(getattr(action, "auto_off_at", None))
def schedule_auto_off(self, action):
seconds = int(action.auto_off_seconds or 0)
if seconds <= 0:
self.cancel_auto_off(action.id)
return
self.cancel_auto_off(action.id)
with self.action_lock:
action.auto_off_at = time.time() + seconds
auto_off_at = action.auto_off_at
timer = threading.Timer(seconds, self.auto_off_action, args=(action.id,))
timer.daemon = True
with self.action_lock:
self.auto_off_timers[action.id] = timer
timer.start()
LOG.info("scheduled auto-off for %s in %ss at %s", action.id, seconds, int(auto_off_at))
def cancel_auto_off(self, action_id):
with self.action_lock:
timer = self.auto_off_timers.pop(action_id, None)
if timer is not None:
timer.cancel()
action = self.action_by_id(action_id)
if action is not None:
with self.action_lock:
action.auto_off_at = None
def auto_off_action(self, action_id):
with self.action_lock:
self.auto_off_timers.pop(action_id, None)
action = self.actions.get(action_id)
if action is None:
return
command = action.safe_off_command()
if not command:
return
try:
LOG.info("auto-off executing %s %s", action_id, command)
action.run(command, bridge=self)
action.auto_off_at = None
self.publish_result(action, command, {"ok": True, "auto_off": True})
except Exception as exc:
LOG.error("auto-off failed for %s: %s", action_id, exc)
self.publish_error(action_id, command, exc)
def publish_result(self, action, command, result):
base = f"{self.base_topic}/{action.id}"
self.publish_state(action, retain=True)
self.mqtt.publish(
f"{base}/result",
json.dumps({"ok": True, "command": command, "time": int(time.time())}),
retain=False,
)
def publish_error(self, action_id, command, exc):
topic = f"{self.base_topic}/{action_id}/result"
payload = {"ok": False, "command": command, "error": str(exc), "time": int(time.time())}
try:
self.mqtt.publish(topic, json.dumps(payload), retain=False)
except Exception:
pass
def all_off(self):
results = []
for action in self.actions_snapshot():
command = action.safe_off_command()
if not command:
continue
try:
action.run(command, bridge=self)
self.cancel_auto_off(action.id)
results.append({"action": action.id, "ok": True})
if action.platform == "switch":
try:
self.publish_state(action, retain=True)
except Exception:
pass
except Exception as exc:
results.append({"action": action.id, "ok": False, "error": str(exc)})
LOG.error("all_off failed for %s: %s", action.id, exc)
return {"ok": all(item["ok"] for item in results), "results": results}
def startup_all_off(self):
if not self.config.get("defaults", {}).get("startup_all_off", True):
return
LOG.info("running startup all-off")
result = self.all_off()
if not result["ok"]:
LOG.warning("startup all-off completed with errors")
def refresh_frigate_inventory(self):
LOG.info("refreshing Frigate camera inventory")
with self.action_lock:
old_actions = dict(self.actions)
new_actions = self.build_actions()
with self.action_lock:
removed_ids = set(old_actions) - set(new_actions)
for action_id in removed_ids:
timer = self.auto_off_timers.pop(action_id, None)
if timer is not None:
timer.cancel()
self.inherit_action_runtime_state(old_actions, new_actions)
self.actions = new_actions
self.refresh_snapshots()
old_ids = set(old_actions)
new_ids = set(new_actions)
old_cameras = self.camera_keys_from_actions(old_actions.values())
new_cameras = self.camera_keys_from_actions(new_actions.values())
result = {
"ok": True,
"added_actions": len(new_ids - old_ids),
"removed_actions": len(old_ids - new_ids),
"added_cameras": len(new_cameras - old_cameras),
"removed_cameras": len(old_cameras - new_cameras),
"total_actions": len(new_ids),
"total_cameras": len(new_cameras),
"mqtt_published": True,
}
result["changed"] = any(
result[key] > 0
for key in ("added_actions", "removed_actions", "added_cameras", "removed_cameras")
)
try:
self.publish_discovery()
except Exception as exc:
result["mqtt_published"] = False
LOG.warning("Frigate refresh could not publish MQTT discovery: %s", exc)
return result
@staticmethod
def camera_keys_from_actions(actions):
keys = set()
for action in actions:
if action.kind == "all_off":
continue
snapshot = action.cfg.get("snapshot") or {}
camera = snapshot.get("camera") or action.cfg.get("camera_name")
if camera:
keys.add((snapshot.get("frigate_url"), camera))
return keys
@staticmethod
def refresh_message(lang, result):
if result.get("added_cameras") or result.get("removed_cameras"):
return text_for(
lang,
"refresh.cameras_changed",
added=result.get("added_cameras", 0),
removed=result.get("removed_cameras", 0),
)
if result.get("added_actions") or result.get("removed_actions"):
return text_for(
lang,
"refresh.controls_changed",
added=result.get("added_actions", 0),
removed=result.get("removed_actions", 0),
)
return text_for(lang, "refresh.no_changes")
def refresh_payload(self, lang, result, redirect_url):
payload = dict(result)
payload["message"] = self.refresh_message(lang, result)
payload["redirect_url"] = redirect_url
return payload
def frigate_refresh_loop(self):
interval = int(self.config.get("defaults", {}).get("frigate_refresh_seconds", 300))
if interval <= 0:
LOG.info("Frigate periodic refresh disabled")
return
interval = max(60, interval)
while not self.stop_event.wait(interval):
try:
self.refresh_frigate_inventory()
except Exception as exc:
LOG.warning("Frigate periodic refresh failed: %s", exc)
def handle_mqtt_message(self, topic, payload):
discovery_prefix = self.config["mqtt"].get("discovery_prefix", "homeassistant")
if topic == f"{discovery_prefix}/status" and payload.lower() == "online":
self.publish_discovery()
return
prefix = f"{self.base_topic}/"
if not topic.startswith(prefix) or not topic.endswith("/set"):
return
action_id = topic[len(prefix) : -len("/set")]
action = self.action_by_id(action_id)
if not action:
LOG.warning("MQTT command for unknown action %s", action_id)
return
value = payload.upper()
if action.platform == "button":
command = "press"
elif value == "ON":
command = "on"
elif value == "OFF":
command = "off"
else:
LOG.warning("unsupported MQTT payload for %s: %s", action_id, payload)
return
try:
LOG.info("MQTT command %s %s", action_id, command)
self.execute_action(action_id, command, publish=True)
except Exception as exc:
LOG.error("MQTT action failed %s %s: %s", action_id, command, exc)
self.publish_error(action_id, command, exc)
def serve(self):
signal.signal(signal.SIGTERM, self._handle_signal)
signal.signal(signal.SIGINT, self._handle_signal)
self.startup_all_off()
mqtt_thread = threading.Thread(
target=self.mqtt.serve,
args=(self.credentials, self.publish_discovery, self.stop_event),
daemon=True,
)
mqtt_thread.start()
snapshot_thread = None
if self.config.get("defaults", {}).get("snapshot_enabled", True):
snapshot_thread = threading.Thread(target=self.snapshot_refresh_loop, daemon=True)
snapshot_thread.start()
refresh_thread = threading.Thread(target=self.frigate_refresh_loop, daemon=True)
refresh_thread.start()
http_cfg = self.config.get("http", {})
if http_cfg.get("enabled", True):
self.start_http(http_cfg.get("bind", "0.0.0.0"), int(http_cfg.get("port", 5011)))
LOG.info("serving %d actions", self.action_count())
while not self.stop_event.is_set():
try:
topic, payload, packet_id = self.messages.get(timeout=1)
self.handle_mqtt_message(topic, payload)
except queue.Empty:
pass
self.mqtt.disconnect_cleanly()
if self.http_server is not None:
self.http_server.shutdown()
mqtt_thread.join(timeout=2)
if snapshot_thread is not None:
snapshot_thread.join(timeout=2)
refresh_thread.join(timeout=2)
def snapshot_refresh_loop(self):
retry_interval = max(60, int(self.config.get("defaults", {}).get("snapshot_retry_seconds", 600)))
refresh_interval = max(retry_interval, int(self.config.get("defaults", {}).get("snapshot_refresh_seconds", 3600)))
interval = retry_interval if self.missing_snapshots() else refresh_interval
while not self.stop_event.wait(interval):
all_ok = self.refresh_snapshots()
interval = refresh_interval if all_ok and not self.missing_snapshots() else retry_interval
def missing_snapshots(self):
for action in self.actions_snapshot():
url = self.snapshot_url(action)
if url and url not in self.snapshot_cache:
return True
return False
def http_token(self):
return self.config.get("http", {}).get("token")
def http_authorized(self, headers, query):
token = self.http_token()
if not token:
return True
values = urllib.parse.parse_qs(query)
supplied = values.get("token", [None])[0] or headers.get("X-Bridge-Token")
return supplied == token
def start_http(self, bind, port):
bridge = self
class Handler(BaseHTTPRequestHandler):
def do_GET(self):
parsed = urllib.parse.urlsplit(self.path)
lang = bridge.request_language(parsed.query, self.headers)
if parsed.path in ("/healthz", "/readyz"):
self.send_json({"ok": True, "actions": bridge.action_count()})
return
if parsed.path.startswith("/static/"):
name = parsed.path.rsplit("/", 1)[-1]
payload = read_static_file(name)
if payload is None:
self.send_response(404)
self.end_headers()
return
content_type = "text/css" if name.endswith(".css") else "application/javascript"
self.send_bytes(payload, content_type, cache_seconds=3600)
return
if not bridge.http_authorized(self.headers, parsed.query):
self.send_response(401)
self.end_headers()
return
if parsed.path.startswith("/snapshot/"):
action_id = urllib.parse.unquote(parsed.path.rsplit("/", 1)[-1]).removesuffix(".jpg")
try:
values = urllib.parse.parse_qs(parsed.query)
refresh = values.get("refresh", ["0"])[0].lower() in ("1", "true", "yes", "on")
live = refresh or bool(bridge.config.get("defaults", {}).get("snapshot_live_on_gui", True))
snapshot = bridge.snapshot_for_action(action_id, live=live)
except Exception as exc:
LOG.info("snapshot request failed for %s: %s", action_id, exc)
snapshot = None
if not snapshot:
self.send_response(404)
self.end_headers()
return
self.send_bytes(snapshot["data"], snapshot["content_type"])
return
if parsed.path == "/api/actions":
self.send_json(bridge.api_actions())
return
token_param = urllib.parse.parse_qs(parsed.query).get("token", [None])[0]
self.send_html(bridge.html_page(lang, token_param))
def do_POST(self):
parsed = urllib.parse.urlsplit(self.path)
lang = bridge.request_language(parsed.query, self.headers)
if not bridge.http_authorized(self.headers, parsed.query):
self.send_response(401)
self.end_headers()
return
if parsed.path == "/refresh":
token_param = urllib.parse.parse_qs(parsed.query).get("token", [None])[0]
token_query = f"&token={urllib.parse.quote(token_param)}" if token_param else ""
redirect_url = f"/?lang={urllib.parse.quote(lang)}{token_query}"
wants_json = "application/json" in self.headers.get("Accept", "").lower()
try:
result = bridge.refresh_frigate_inventory()
if wants_json:
self.send_json(bridge.refresh_payload(lang, result, redirect_url))
else:
self.send_response(303)
self.send_header("Location", redirect_url)
self.end_headers()
except Exception as exc:
if wants_json:
self.send_json(
{
"ok": False,
"changed": False,
"message": text_for(lang, "refresh.failed", error=str(exc)),
"redirect_url": redirect_url,
},
status=500,
)
else:
self.send_response(500)
self.end_headers()
self.wfile.write(str(exc).encode("utf-8"))
return
parts = parsed.path.strip("/").split("/")
if len(parts) == 3 and parts[0] == "action":
action_id = urllib.parse.unquote(parts[1])
command = urllib.parse.unquote(parts[2]).lower()
try:
result = bridge.execute_action(action_id, command, publish=True)
token_param = urllib.parse.parse_qs(parsed.query).get("token", [None])[0]
token_query = f"&token={urllib.parse.quote(token_param)}" if token_param else ""
self.send_response(303)
self.send_header("Location", f"/?lang={urllib.parse.quote(lang)}{token_query}")
self.end_headers()
except Exception as exc:
self.send_response(500)
self.end_headers()
self.wfile.write(str(exc).encode("utf-8"))
return
self.send_response(404)
self.end_headers()
def send_json(self, payload, status=200):
data = json.dumps(payload, indent=2).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(data)))
self.end_headers()
self.wfile.write(data)
def send_html(self, body):
data = body.encode("utf-8")
self.send_response(200)
self.send_header("Content-Type", "text/html; charset=utf-8")
self.send_header("Content-Length", str(len(data)))
self.end_headers()
self.wfile.write(data)
def send_bytes(self, data, content_type, cache_seconds=300):
self.send_response(200)
self.send_header("Content-Type", content_type)
self.send_header("Cache-Control", f"private, max-age={cache_seconds}")
self.send_header("Content-Length", str(len(data)))
self.end_headers()
self.wfile.write(data)
def log_message(self, fmt, *args):
LOG.info("http %s - %s", self.address_string(), fmt % args)
self.http_server = ThreadingHTTPServer((bind, port), Handler)
thread = threading.Thread(target=self.http_server.serve_forever, daemon=True)
thread.start()
LOG.info("HTTP fallback GUI listening on %s:%s", bind, port)
def api_actions(self):
return [
{
"id": action.id,
"name": action.display_name(self.default_language),
"platform": action.platform,
"commands": action.command_list(),
"state": action.state,
"snapshot": bool(action.cfg.get("snapshot")),
"auto_off_seconds": action.auto_off_seconds,
"auto_off_at": int(action.auto_off_at) if action.auto_off_at else None,
"requires_ack": bool(action.cfg.get("requires_ack")),
}
for action in self.actions_snapshot()
]
def html_page(self, lang=None, token_param=None):
with self.action_lock:
actions = list(self.actions.values())
return render_page(actions, lang or DEFAULT_LANGUAGE, token_param)
def request_language(self, query, headers=None):
values = urllib.parse.parse_qs(query)
requested = values.get("lang", [None])[0]
if requested:
return normalize_language(requested)
accept_language = headers.get("Accept-Language") if headers else None
return language_from_accept_header(accept_language)
def _handle_signal(self, signum, _frame):
LOG.info("received signal %s, stopping", signum)
self.stop_event.set()