135 lines
5.4 KiB
Python
135 lines
5.4 KiB
Python
import asyncio
|
|
import time
|
|
import requests
|
|
import datetime
|
|
from typing import Dict
|
|
from zeroconf import ServiceStateChange
|
|
from zeroconf.asyncio import AsyncZeroconf, AsyncServiceBrowser, AsyncServiceInfo
|
|
|
|
SERVICE_TYPE = "_displaytcp._tcp.local."
|
|
|
|
class DisplayManager:
|
|
def __init__(self, endpoint_path: str = "/api/display"):
|
|
self.endpoint_path = endpoint_path
|
|
self.screens: Dict[str, str] = {} # Map: service_name -> URL
|
|
self.aiozc: AsyncZeroconf | None = None
|
|
self.browser: AsyncServiceBrowser | None = None
|
|
self._on_screens_change_callback = None # Callback for screen changes
|
|
|
|
async def start(self):
|
|
"""Starts dynamic mDNS discovery for e-Paper screens."""
|
|
self.aiozc = AsyncZeroconf()
|
|
self.browser = AsyncServiceBrowser(
|
|
self.aiozc.zeroconf,
|
|
SERVICE_TYPE,
|
|
handlers=[self._on_service_state_change]
|
|
)
|
|
print(f"[DisplayManager] 🔍 Listening for mDNS screens ({SERVICE_TYPE})...")
|
|
|
|
def _on_service_state_change(self, zeroconf, service_type, name, state_change):
|
|
asyncio.create_task(self._update_service(name, state_change))
|
|
|
|
async def _update_service(self, name: str, state_change: ServiceStateChange):
|
|
if state_change in (ServiceStateChange.Added, ServiceStateChange.Updated):
|
|
info = AsyncServiceInfo(SERVICE_TYPE, name)
|
|
if await info.async_request(self.aiozc.zeroconf, 3000):
|
|
addresses = info.parsed_addresses()
|
|
if addresses:
|
|
# Prefer IPv4 if available, fallback to IPv6 (Thread)
|
|
ipv4_addrs = [a for a in addresses if ":" not in a]
|
|
|
|
if ipv4_addrs:
|
|
ip_str = ipv4_addrs[0]
|
|
else:
|
|
raw_ip = addresses[0]
|
|
# Wrap IPv6 addresses in square brackets for valid HTTP URLs
|
|
ip_str = f"[{raw_ip}]" if not raw_ip.startswith("[") else raw_ip
|
|
|
|
port = info.port or 5000
|
|
url = f"http://{ip_str}:{port}{self.endpoint_path}"
|
|
|
|
self.screens[name] = url
|
|
if self._on_screens_change_callback:
|
|
self._on_screens_change_callback(self.screens)
|
|
print(f"[DisplayManager] 📺 Screen registered: {name} -> {url}")
|
|
|
|
elif state_change == ServiceStateChange.Removed:
|
|
if name in self.screens:
|
|
print(f"[DisplayManager] ❌ Screen disconnected: {name}")
|
|
self.screens.pop(name, None)
|
|
if self._on_screens_change_callback:
|
|
self._on_screens_change_callback(self.screens)
|
|
|
|
def set_on_screens_change_callback(self, callback):
|
|
"""Set a callback function to be called when screens are added or removed."""
|
|
self._on_screens_change_callback = callback
|
|
|
|
async def broadcast_state(self, microwave_states: dict, cloud_alert: bool):
|
|
"""Broadcasts the system state JSON to all discovered screens concurrently."""
|
|
if not self.screens:
|
|
return
|
|
|
|
now_time = time.time()
|
|
update_time_str = datetime.datetime.now().strftime("%H:%M")
|
|
|
|
microwaves_list = []
|
|
for mw_id, state_info in microwave_states.items():
|
|
# Extract state properties safely
|
|
state_val = state_info.cooking_state
|
|
paused = state_info.paused
|
|
rem_time = state_info.cooking_estimated_remaining_time
|
|
start_t = state_info.cooking_start_time
|
|
name = state_info.friendly_name
|
|
|
|
# Calculate progress percentage (0 to 100)
|
|
progress = 0
|
|
if start_t > 0 and rem_time > 0:
|
|
elapsed = now_time - start_t
|
|
total = elapsed + rem_time
|
|
if total > 0:
|
|
progress = int((elapsed / total) * 100)
|
|
progress = max(0, min(100, progress))
|
|
|
|
microwaves_list.append({
|
|
"id": mw_id,
|
|
"name": name,
|
|
"state": state_val,
|
|
"paused": paused,
|
|
"remaining_time": rem_time,
|
|
"progress": progress
|
|
})
|
|
|
|
payload = {
|
|
"cloud_alert": cloud_alert,
|
|
"update_time": update_time_str,
|
|
"microwaves": microwaves_list
|
|
}
|
|
|
|
# Dispatch POST requests in parallel without blocking the main loop
|
|
tasks = [
|
|
self._post_to_screen(name, url, payload)
|
|
for name, url in list(self.screens.items())
|
|
]
|
|
await asyncio.gather(*tasks, return_exceptions=True)
|
|
|
|
async def _post_to_screen(self, name: str, url: str, payload: dict):
|
|
try:
|
|
response = await asyncio.to_thread(
|
|
requests.post,
|
|
url,
|
|
json=payload,
|
|
headers={"Content-Type": "application/json"},
|
|
timeout=3.0
|
|
)
|
|
response.raise_for_status()
|
|
print(f"[DisplayManager] Updated {name}")
|
|
except Exception as e:
|
|
print(f"[DisplayManager] Error pushing state to {name}: {e}")
|
|
|
|
async def stop(self):
|
|
"""Clean up mDNS browser and Zeroconf instance."""
|
|
if self.browser:
|
|
await self.browser.async_cancel()
|
|
if self.aiozc:
|
|
await self.aiozc.async_close()
|
|
print("[DisplayManager] Stopped service browser.") |