"""Async event bus with pluggable sinks (log, webhook, email). Events are published from the pipeline thread and delivered on a dedicated worker thread, so a slow webhook or SMTP server can never stall frame processing. Sink failures are logged, never raised. """ from __future__ import annotations import logging import queue import smtplib import threading import time from collections import deque from dataclasses import asdict, dataclass, field from email.mime.text import MIMEText log = logging.getLogger(__name__) @dataclass class Event: type: str # person.new | person.seen | camera.up | camera.down | ... camera_id: str ts: float = field(default_factory=time.time) data: dict = field(default_factory=dict) def to_dict(self) -> dict: return asdict(self) class EventBus: def __init__(self) -> None: self._queue: "queue.Queue[Event | None]" = queue.Queue(maxsize=1000) self._sinks: list = [] self.recent: deque = deque(maxlen=300) self._worker = threading.Thread( target=self._run, daemon=True, name="event-bus") self._started = False def add_sink(self, sink) -> None: self._sinks.append(sink) def start(self) -> None: if not self._started: self._started = True self._worker.start() def stop(self) -> None: if self._started: self._queue.put(None) self._worker.join(timeout=5) def publish(self, event: Event) -> None: self.recent.appendleft(event.to_dict()) try: self._queue.put_nowait(event) except queue.Full: log.warning("event queue full, dropping %s", event.type) def _run(self) -> None: while True: event = self._queue.get() if event is None: return for sink in self._sinks: try: sink.handle(event) except Exception: log.exception("sink %s failed for %s", type(sink).__name__, event.type) class LogSink: def handle(self, event: Event) -> None: log.info("event %s [%s] %s", event.type, event.camera_id, event.data) class WebhookSink: def __init__(self, url: str, timeout: float = 5.0): self.url = url self.timeout = timeout def handle(self, event: Event) -> None: import requests requests.post(self.url, json=event.to_dict(), timeout=self.timeout) class EmailSink: """Rate-limited email notifications for person events only.""" NOTIFY_TYPES = {"person.new", "person.seen"} def __init__(self, smtp_host: str, smtp_port: int, username: str, password: str, to: str, min_interval: float = 300.0): self.smtp_host = smtp_host self.smtp_port = smtp_port self.username = username self.password = password self.to = to self.min_interval = min_interval self._last_sent = 0.0 def handle(self, event: Event) -> None: if event.type not in self.NOTIFY_TYPES: return now = time.time() if now - self._last_sent < self.min_interval: return self._last_sent = now label = event.data.get("label", "someone") body = (f"Behavision: {label} detected on camera {event.camera_id}\n" f"Event: {event.type}\nDetails: {event.data}") msg = MIMEText(body) msg["Subject"] = f"Behavision: {label} on {event.camera_id}" msg["From"] = self.username msg["To"] = self.to with smtplib.SMTP(self.smtp_host, self.smtp_port, timeout=10) as smtp: smtp.starttls() smtp.login(self.username, self.password) smtp.send_message(msg)