doorbell-listener: flash plug on alert (configurable interval/duration)
On a fresh (< MAX_AGE_SECONDS) ntfy message, spawn a thread that publishes N TOGGLEs spaced FLASH_INTERVAL_SECONDS apart so total runtime is FLASH_DURATION_SECONDS. Defaults: 5 toggles @ 3s = 15s. Env vars (all in compose.yaml): FLASH_INTERVAL_SECONDS default 3 FLASH_DURATION_SECONDS default 15 MAX_AGE_SECONDS default 180 (replay protection) Each alert gets its own thread; concurrent alerts spawn concurrent flashes. z2m just flips state per toggle, so the lamp alternates regardless of overlap.
This commit is contained in:
@@ -57,6 +57,8 @@ services:
|
|||||||
# Listens to ntfy.sh for doorbell events and toggles a smart plug on each alert.
|
# Listens to ntfy.sh for doorbell events and toggles a smart plug on each alert.
|
||||||
# Topic: ALERT_klubhaus_topic_test (matches ESP32 firmware DEBUG_MODE suffix)
|
# Topic: ALERT_klubhaus_topic_test (matches ESP32 firmware DEBUG_MODE suffix)
|
||||||
# Action: TOGGLE on zigbee2mqtt/Sideboard Lamp/set
|
# Action: TOGGLE on zigbee2mqtt/Sideboard Lamp/set
|
||||||
|
# On a fresh alert, flash the plug: N toggles spaced FLASH_INTERVAL_SECONDS
|
||||||
|
# so the total runtime is FLASH_DURATION_SECONDS.
|
||||||
doorbell-listener:
|
doorbell-listener:
|
||||||
build: ./doorbell-listener
|
build: ./doorbell-listener
|
||||||
container_name: doorbell-listener
|
container_name: doorbell-listener
|
||||||
@@ -70,6 +72,9 @@ services:
|
|||||||
MQTT_PORT: "1883"
|
MQTT_PORT: "1883"
|
||||||
PLUG_TOPIC: "zigbee2mqtt/Sideboard Lamp/set"
|
PLUG_TOPIC: "zigbee2mqtt/Sideboard Lamp/set"
|
||||||
POLL_INTERVAL: "30"
|
POLL_INTERVAL: "30"
|
||||||
|
MAX_AGE_SECONDS: "180"
|
||||||
|
FLASH_INTERVAL_SECONDS: "3"
|
||||||
|
FLASH_DURATION_SECONDS: "15"
|
||||||
TZ: America/Los_Angeles
|
TZ: America/Los_Angeles
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ Env vars (all optional):
|
|||||||
import json
|
import json
|
||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
|
import threading
|
||||||
import time
|
import time
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
@@ -44,6 +45,13 @@ MQTT_PORT = int(os.environ.get("MQTT_PORT", "1883"))
|
|||||||
PLUG_TOPIC = os.environ.get("PLUG_TOPIC", "zigbee2mqtt/Sideboard Lamp/set")
|
PLUG_TOPIC = os.environ.get("PLUG_TOPIC", "zigbee2mqtt/Sideboard Lamp/set")
|
||||||
STATE_DIR = Path(os.environ.get("STATE_DIR", "/data/doorbell-listener"))
|
STATE_DIR = Path(os.environ.get("STATE_DIR", "/data/doorbell-listener"))
|
||||||
POLL_INTERVAL = int(os.environ.get("POLL_INTERVAL", "30"))
|
POLL_INTERVAL = int(os.environ.get("POLL_INTERVAL", "30"))
|
||||||
|
# Drop messages older than this many seconds (replay protection after restart).
|
||||||
|
MAX_AGE_SECONDS = int(os.environ.get("MAX_AGE_SECONDS", "180"))
|
||||||
|
# On a fresh alert, flash the plug: TOGGLE every FLASH_INTERVAL_SECONDS,
|
||||||
|
# repeated so the total runtime is FLASH_DURATION_SECONDS. With defaults
|
||||||
|
# (3s interval, 15s duration) that's 5 toggles, ending opposite of start.
|
||||||
|
FLASH_INTERVAL_SECONDS = float(os.environ.get("FLASH_INTERVAL_SECONDS", "3"))
|
||||||
|
FLASH_DURATION_SECONDS = float(os.environ.get("FLASH_DURATION_SECONDS", "15"))
|
||||||
|
|
||||||
|
|
||||||
def log(msg: str) -> None:
|
def log(msg: str) -> None:
|
||||||
@@ -67,6 +75,26 @@ def save_last_id(topic: str, msg_id: str) -> None:
|
|||||||
last_id_path(topic).write_text(msg_id)
|
last_id_path(topic).write_text(msg_id)
|
||||||
|
|
||||||
|
|
||||||
|
def start_flash(client: mqtt.Client, count: int, interval: float) -> None:
|
||||||
|
"""Spawn a thread that publishes `count` TOGGLEs spaced `interval` seconds.
|
||||||
|
|
||||||
|
Each alert gets its own thread; overlapping flashes will publish toggles
|
||||||
|
concurrently but z2m just flips state per toggle, so the lamp alternates
|
||||||
|
regardless of how many threads are running.
|
||||||
|
"""
|
||||||
|
def run() -> None:
|
||||||
|
log(f"flash: starting {count} toggles @ {interval}s")
|
||||||
|
for i in range(count):
|
||||||
|
payload = json.dumps({"state": "TOGGLE"})
|
||||||
|
client.publish(PLUG_TOPIC, payload)
|
||||||
|
log(f"flash {i + 1}/{count} -> {PLUG_TOPIC}")
|
||||||
|
if i < count - 1:
|
||||||
|
time.sleep(interval)
|
||||||
|
log("flash: complete")
|
||||||
|
|
||||||
|
threading.Thread(target=run, daemon=True).start()
|
||||||
|
|
||||||
|
|
||||||
def poll_topic(client: mqtt.Client, topic: str, last_id: str) -> str:
|
def poll_topic(client: mqtt.Client, topic: str, last_id: str) -> str:
|
||||||
"""Poll one topic. Returns the latest message id seen (or last_id)."""
|
"""Poll one topic. Returns the latest message id seen (or last_id)."""
|
||||||
url = f"https://ntfy.sh/{topic}/json?poll=1"
|
url = f"https://ntfy.sh/{topic}/json?poll=1"
|
||||||
@@ -78,6 +106,8 @@ def poll_topic(client: mqtt.Client, topic: str, last_id: str) -> str:
|
|||||||
r.raise_for_status()
|
r.raise_for_status()
|
||||||
|
|
||||||
newest_id = last_id
|
newest_id = last_id
|
||||||
|
now = time.time()
|
||||||
|
flash_count = max(1, int(FLASH_DURATION_SECONDS / FLASH_INTERVAL_SECONDS))
|
||||||
for line in r.iter_lines(decode_unicode=True):
|
for line in r.iter_lines(decode_unicode=True):
|
||||||
if not line:
|
if not line:
|
||||||
continue
|
continue
|
||||||
@@ -89,18 +119,25 @@ def poll_topic(client: mqtt.Client, topic: str, last_id: str) -> str:
|
|||||||
continue
|
continue
|
||||||
|
|
||||||
msg_id = event.get("id", "")
|
msg_id = event.get("id", "")
|
||||||
|
msg_time = event.get("time", 0)
|
||||||
title = event.get("title", "")
|
title = event.get("title", "")
|
||||||
message = event.get("message", "")
|
message = event.get("message", "")
|
||||||
log(f"[{topic}] alert: id={msg_id} title={title!r} message={message!r}")
|
|
||||||
|
|
||||||
payload = json.dumps({"state": "TOGGLE"})
|
|
||||||
info = client.publish(PLUG_TOPIC, payload)
|
|
||||||
info.wait_for_publish(timeout=5)
|
|
||||||
log(f"published TOGGLE on {PLUG_TOPIC}")
|
|
||||||
|
|
||||||
|
# Update last_id for every message we see, even stale ones, so
|
||||||
|
# restarts don't replay them.
|
||||||
if msg_id:
|
if msg_id:
|
||||||
newest_id = msg_id
|
newest_id = msg_id
|
||||||
save_last_id(topic, msg_id)
|
save_last_id(topic, msg_id)
|
||||||
|
|
||||||
|
# Freshness filter: only act on recent messages.
|
||||||
|
age = int(now - msg_time) if msg_time else None
|
||||||
|
if age is not None and age > MAX_AGE_SECONDS:
|
||||||
|
log(f"[{topic}] skip stale: id={msg_id} age={age}s > {MAX_AGE_SECONDS}s")
|
||||||
|
continue
|
||||||
|
|
||||||
|
log(f"[{topic}] alert: id={msg_id} age={age}s title={title!r} message={message!r}")
|
||||||
|
start_flash(client, flash_count, FLASH_INTERVAL_SECONDS)
|
||||||
|
|
||||||
return newest_id
|
return newest_id
|
||||||
|
|
||||||
|
|
||||||
@@ -110,6 +147,9 @@ def main() -> None:
|
|||||||
log(f" mqtt broker: {MQTT_HOST}:{MQTT_PORT}")
|
log(f" mqtt broker: {MQTT_HOST}:{MQTT_PORT}")
|
||||||
log(f" plug topic: {PLUG_TOPIC}")
|
log(f" plug topic: {PLUG_TOPIC}")
|
||||||
log(f" poll interval: {POLL_INTERVAL}s")
|
log(f" poll interval: {POLL_INTERVAL}s")
|
||||||
|
log(f" max age: {MAX_AGE_SECONDS}s")
|
||||||
|
log(f" flash: {FLASH_DURATION_SECONDS / FLASH_INTERVAL_SECONDS:.0f} toggles "
|
||||||
|
f"@ {FLASH_INTERVAL_SECONDS}s ({FLASH_DURATION_SECONDS}s total)")
|
||||||
log(f" state dir: {STATE_DIR}")
|
log(f" state dir: {STATE_DIR}")
|
||||||
|
|
||||||
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2)
|
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2)
|
||||||
|
|||||||
Reference in New Issue
Block a user