Files
LEDMatrix/src/common/sync_manager.py
T
ChuckandClaude Opus 5.5 7b90759252 fix: /errors stack traces, Wi-Fi disconnect and save, plugin fonts, API cache TTL (#636)
* fix(errors): record the exception's own stack trace

record_error() called traceback.format_exc(), which only sees an
exception while its except block is running. plugin_executor records
exceptions caught on a worker thread after that block has ended, so
every trace on /errors read "NoneType: None". The trace is now built
from the exception's __traceback__. The executor's log call had the
same problem with exc_info=True and now passes the exception.

record_error() also merged LEDMatrixError context into the caller's
dict in place; it now works on a copy.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* docs(wifi): point at configure_wifi_permissions.sh instead of a sudoers list

The module docstring told users to grant NOPASSWD sudo on iptables and
ip. configure_wifi_permissions.sh refuses those grants on purpose: a
wildcard rule for either runs an arbitrary program as root. Point at
the script and say why it leaves them out.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(wifi): disconnect finds the saved profile by SSID

disconnect_from_network() asked `nmcli -f NAME,802-11-wireless.ssid
connection show` for the profile to take down, but nmcli rejects that
column for `connection show`, so the lookup always failed and only the
device was disconnected. The per-profile lookup _connect_nmcli() already
used is now _find_profile_for_ssid(), and both callers share it. It
also splits terse output on the last colon and unescapes "\:", so a
profile name containing a colon is found.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(wifi): write wifi_config.json atomically and report a failed save

_save_config() opened the file for writing in place and swallowed any
error, so a wifi_config.json left owned by root made the web toggle for
auto-enabling AP mode report success while nothing was saved, and a
crash mid-write could truncate the file. It now uses atomic_write_json,
which also keeps the file's owner and shared group when root saves it,
and returns False on failure. POST /wifi/ap/auto-enable answers 500 in
that case.

The file is now written with indent=4, like the other config files.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(fonts): resolve plugin:// fonts in the plugin's own directory

FontManager looked for a plugin's bundled fonts under Path("plugins") /
plugin_id: relative to the process cwd, and not the default install
directory (plugin-repos/), so a manifest's plugin:// fonts never loaded.

register_plugin_fonts() takes an optional plugin_dir, and PluginManager
passes the directory it loaded the plugin from. Callers that omit it get
a lookup in the configured plugin_system.plugins_directory, then plugins/,
resolved against the install root.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(api-helper): cache responses for the requested cache_ttl

APIHelper.get(cache_ttl=...) and set_cache(ttl=...) dropped the ttl on
the claim that CacheManager does not support one, but CacheManager.set()
takes a ttl, stores it with the entry, and both cache tiers honour it
over a reader's max_age. Without it every response expired after the
300-second default read age, whatever the plugin asked for. The ttl is
now passed through, and the cache read passes cache_ttl as max_age for
entries written without one. The class docstring describes what the
helper actually does.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(style): one scale range for the schema, element_scale and LogoHelper

The generated Scale field allowed 0.1 to 10, element_style's reader
capped at 10 with no floor, and LogoHelper accepted 0.05 to 8 and reset
anything else to 1.0. A logo scale of 9, which the form accepts, drew at
the shipped size.

MIN_ELEMENT_SCALE / MAX_ELEMENT_SCALE (0.1, 10.0) in src.element_style
are now the schema bounds and the clamp every reader applies through
coerce_scale(): a positive number outside the range is clamped, and
anything that is not a finite positive number means the default. That
also stops element_scale() passing NaN through, since min(nan, 10.0)
is nan.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(logos): placeholder lands at the requested path; empty logos list

download_missing_logo() wrote its fallback placeholder to
<normalize_abbreviation(abbr)>.png in the logo directory rather than to
the logo_path the caller passed, so it could return True while nothing
existed where the plugin looks (e.g. "TA&M.png" vs "TAANDM.png").
create_placeholder_logo() takes an optional filepath, and
download_missing_logo passes the requested one.

download_missing_logo_for_team() only caught KeyError, so a team whose
"logos" list is empty raised IndexError; it now treats KeyError,
IndexError and TypeError as "no logo URL".

The placeholder is drawn with PLACEHOLDER_SIZE / PLACEHOLDER_BG, the
constants is_placeholder_logo() recognises it by, instead of repeated
literals.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(fonts): resolve bundled font paths against the install root

TextHelper's default font_dir, the logo placeholder's font and
FontManager's font_overrides.json were all relative to the process cwd,
so a process started anywhere but the install root (the plugin safety
harness, a manual run, a unit without WorkingDirectory) drew with PIL's
default face and read no overrides. They now go through
font_layout.resolve_asset_path; the overrides file sits in the install
root's config/.

The resolver docstrings described an order the code does not follow:
resolve_asset_path never consults the cwd, and sports_shared's
_resolve_font_path tries the cwd first. Both docstrings now say what
the code does, and _resolve_font_path calls resolve_asset_path instead
of probing FontManager for it.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(sync): the web UI reads the sync status file the display writes

sync_manager writes its status to tempfile.gettempdir(), but
GET /api/v3/sync/status read a hardcoded /tmp/led_matrix_sync_status.json
and defaulted the port to a literal 5765. Wherever TMPDIR is set (or on
any non-/tmp host) the page only ever showed "starting". The endpoint now
uses sync_manager.STATUS_FILE and SYNC_PORT.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(http): the rankings resolver sends the project's User-Agent

DynamicTeamResolver fetched ESPN rankings with a bare requests.get, so
it sent python-requests' default User-Agent, which ESPN rejects; the
AP_TOP_N favourites then resolved to nothing. It now sends
DEFAULT_HTTP_HEADERS. BaseOddsManager carried its own copy of the
User-Agent string and now uses the same shared headers (which also adds
Accept-Language).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(backup): record the core release and read the configured plugin dir

The manifest's ledmatrix_version came from a VERSION file that does not
exist, then from .git/HEAD: a 12-character sha, or "ref: refs/he" when
the branch's ref was packed. It is now src.__version__.

list_installed_plugins() scanned a hardcoded plugin-repos/, so on an
install whose plugin_system.plugins_directory points elsewhere, plugins
missing from plugin_state.json were left out of the backup. It now reads
the configured directory from config/config.json, defaulting to
plugin-repos.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(startup): report a missing display section once

A config without a display section produced three errors for the one
problem ("Missing required configuration key: display", "Display
configuration is missing or empty" and "Display configuration is
missing"), and an empty one produced two. _validate_config now reports
it once, as a missing key or an empty section, and
_validate_display_config leaves it to that.

The module docstring said the validator fails fast; nothing in the
display service calls raise_on_errors(), so it now says the errors are
reported and startup continues.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(wifi): share the copied blocks and name the AP constants

- _parse_nmcli_wifi_list() is the one parser behind _scan_nmcli and
  _scan_nmcli_cached.
- _verify_connected(), _wait_for_device_idle(), _failsafe_ap() and
  _mark_forced() replace blocks that were pasted two or three times in
  the connect and enable-AP paths. The device-idle wait now checks
  before its first one-second sleep instead of after it.
- _check_command() calls _find_command_path() instead of repeating it.
- AP_IP, PORTAL_PORT, AP_PROFILE_NAME and AP_PROFILE_NAMES name values
  that were spelled out 14, 12, 8 and 2 times; the two deletion loops
  now walk the same tuple. The iwconfig status path compares the AP
  address exactly: startswith() also skipped 192.168.4.10-19.
- Dropped a second WIFI.SIGNAL query that repeated the first, a no-op
  "if ssid: continue", the try/except around _connect_wpa_supplicant's
  constant return, and a second save of a scan scan_networks already
  saves.
- _ensure_wifi_radio_enabled's docstring says it returns True when the
  radio state cannot be read at all.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(config): drop dead branches and history comments in ConfigManager

- The module docstring pointed plugin authors at update_plugin_config(),
  which does not exist; it now names save_config_atomic() and
  save_raw_file_content().
- load_config's FileNotFoundError handler tested the message for
  "config_secrets.json", but a missing secrets file is handled where it
  is read, so only config.json reaches it; the check is gone.
- save_raw_file_content's `file_type == "main" or "secrets"` guard was
  always true (anything else raised earlier).
- get_raw_file_content('secrets') already returns {} for a missing file,
  so the os.path.exists() in front of two calls to it is gone.
- Comments that narrated earlier behaviour are rewritten as what the
  code does now.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(background-data): present-tense comments, drop unused API

- Comments that told the history of each fix (what "used to" happen,
  "the old per-delivery release") now state the invariant the code keeps.
- get_statistics() no longer reports a constant 'queue_size': 0, and the
  uncalled clear_completed_requests() is gone (_cleanup_completed_requests
  does that job on every completion). Neither is referenced in core, the
  web UI or the plugin monorepo.

shutdown_background_service() has no production caller either, but it
is the only way to tear down the get_background_service() singleton,
which the tests rely on, so it stays.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(odds): drop the unread cache_ttl and merge the odds_data branches

BaseOddsManager loaded base_odds_manager.cache_ttl from config and never
used it: cached odds live for the update interval (get_odds' ttl=interval).
No core or monorepo code reads the attribute, so it is gone along with
its log line. The two consecutive `if odds_data:` blocks are one.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(backup): one table for the single-file sections

config, secrets, wifi and ytm_auth were each spelled out in create,
preview, validate and restore. _SINGLE_FILE_SECTIONS lists them once,
with the RestoreOptions flag that restores each, and all four walk it.
Restore error messages keep their wording ("Failed to restore
<file name>").

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(fonts): drop FontManager's write-only state and duplicate logs

- fonts_config, font_metadata and font_dependencies were written and
  never read; the performance_stats keys font_load_times, render_times,
  total_renders and the per-call "resolve" timings
  (_record_performance_metric) likewise. get_performance_stats() reads
  only the counters that remain. Nothing in core or the plugin monorepo
  references any of them.
- A failed BDF load was logged twice, by _load_bdf_font and again by
  get_font; get_font's line is the one kept.
- Removed "NEW:" and commented-out cozette entries, the "Copy font to
  assets/fonts" comment on code that copies nothing, and local imports
  of names the module already imports. The deprecated add_font() now
  resolves assets/fonts against the install root.

The @deprecated methods stay.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(text-helper): cache loaded fonts; drop the pre-textlength fallback

TextHelper declared _font_cache, cleared it and reported its size, but
never stored anything in it. load_fonts() now keeps each (file, size)
it loads there, so clear_font_cache() and get_font_cache_stats() mean
what they say and repeated load_fonts() calls reuse the fonts.

get_text_width() no longer catches AttributeError for Pillow releases
without ImageDraw.textlength; requirements.txt pins Pillow>=12.2.
The class docstring describes what the helper does.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* docs(common): fix wrong docstrings in api_helper, permission_utils, snapshot_policy

- permission_utils called 0o2775 "sticky bit"; the 2 is setgid, which is
  what makes new files take the directory's group.
- snapshot_policy pointed at web_interface/blueprints/api_v3.py, which
  is a package now; the health check is in api_v3/misc.py.
- APIHelper.clear_cache() lost a history note and a fallback to a
  clear() method that neither CacheManager nor the testing
  MockCacheManager has. The session headers are built from
  DEFAULT_HTTP_HEADERS instead of a copy of them, and the module
  docstring says what the module offers.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* docs(sports): present-tense comments in the shared scoreboard renderers

- sports_scroll and sports_game_renderer comments that referred to "this
  PR", "the old flat 128px card" or what the renderer "previously" did
  now describe the current behaviour and its reason.
- The block explaining why non-finite settings are rejected sat above
  _score_reserve_width; it describes _center_gap_width and now lives in
  it.
- unshare_element_fonts wrapped its import of font_layout.load_truetype
  in an `except ImportError` that cannot fire inside core; the import
  stays at call time so tests can spy on the pinned loader.
- sports_card docstrings that told the history of a fix say what the
  code does.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(sports-shared): drop dead code, name the ESPN limit

- _get_weeks_data asked for limit=1000, which fetch_espn_scoreboard
  clamps to ESPN_MAX_LIMIT anyway; it now names that constant. Its
  unused `immediate_events = []` is gone.
- _get_season_schedule_dates() returned ("", "") and has no caller in
  core or the plugin monorepo.
- _should_log keeps its warning_type parameter (part of the inherited
  signature, though nothing in core or the monorepo calls it) and its
  docstring says the cooldown is shared across types.
- An unused ImageFont import is gone.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(sync): one follower-mode switch, shared panel defaults

- The class docstring said the leader sends PNG frames. Frames go over
  UDP as raw RGB; PNG is only the Vegas scroll image sent over TCP. It
  now describes both paths.
- _enter_follower_mode() replaces the two copies of "note the leader,
  switch from standalone to follower, log, write status" in the frame
  and scroll-position handlers.
- The rows/cols fallbacks use DEFAULT_ROWS / DEFAULT_COLS from
  src.display_geometry, as chain_length already did.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(style): drop _layout_axis, name the layout group title

- ElementStyleResolver._layout_axis() had no caller in core or the
  plugin monorepo.
- _element_block_from_spec checked spec['size'] was a dict again after
  size_spec already had; it reads size_spec.
- The "Layout Offsets" title written into three generated schema blocks
  is _LAYOUT_TITLE.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* docs(logo-helper): say what the placeholder draws; name the 1.5 box factor

- _create_placeholder_logo's docstring said it draws the team
  abbreviation; it draws an outlined grey box and nothing else. The
  docstring says so, and the "in a real implementation you'd want text"
  comments are gone.
- The 1.5 x panel default logo box, written out six times, is
  DEFAULT_LOGO_BOX_FACTOR.
- ImageDraw is imported with Image at the top of the module.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(logos): drop dead code and a duplicate regex in logo_downloader

- _SAFE_LEAGUE_CODE_RE was the same pattern as _SAFE_LEAGUE_RE; both
  checks use the one.
- get_logo_filename_variations reassigned the TA&M case to the list it
  already had; the function returns the two names directly.
- _get_team_name_variations() had no caller in core or the plugin
  monorepo.
- fetch_single_team's docstring was copied from fetch_teams_data; a log
  message read "for{team_id}".

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor: drop the Pillow<9.1 resample shim and a catch-and-reraise

- adaptive_images fell back to Image.LANCZOS/NEAREST for Pillow < 9.1;
  requirements.txt pins Pillow>=12.2. RESAMPLE_LANCZOS and
  RESAMPLE_NEAREST keep their names (src.common re-exports them).
- CacheManager.save_cache caught CacheError only to re-raise it; the
  disk write is now called directly, with the same result.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* test(api-helper): stop the real CacheManager's cleanup thread

The cache-lifetime tests built a CacheManager and left its cleanup
thread's class-wide claim on the directory in place, which broke
test_cache_cleanup_thread_ownership when it ran later in the session.
The fixture now stops the thread on teardown.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* docs(changelog): core-common

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-24 17:32:29 -04:00

714 lines
31 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
Multi-Display Sync Manager
Synchronizes scrolling content across two LED matrix display units over UDP.
Runs at the core framework level — works with any plugin automatically.
Roles:
standalone No sync (default behavior)
leader Drives scroll, sends rendered follower frames via UDP
follower Receives frames from leader; falls back to own plugins when
the leader goes offline
Compatibility rule: rows and cols must match between leader and follower.
chain_length may differ — each display can have a different number of panels.
Port default: 5765 (UDP). Open this port on both Pis if ufw is active:
sudo ufw allow 5765/udp
"""
import io
import json
import math
import os
import socket
import struct
import tempfile
import threading
import time
import logging
from enum import Enum
from typing import Callable, Optional
import numpy as np
from PIL import Image
from src.display_geometry import DEFAULT_CHAIN_LENGTH, DEFAULT_COLS, DEFAULT_ROWS
# Raw-frame wire format: 8-byte magic + 4-byte header + raw RGB pixels
# Much faster than PNG: no encode/decode, negligible CPU, same UDP packet size
_RAW_MAGIC = b'SYNC_RAW'
_RAW_HEADER = struct.Struct('<HH') # width, height (uint16 LE)
# Upper bound on a decoded frame/scroll image. Generous for any real scroll
# image (a leader's full cycle is long but only panel-height tall), and low
# enough that a crafted image from any host on the LAN cannot force a large
# allocation on the render thread. Applied on both receive paths — the TCP
# image server and the follower's legacy-PNG UDP fallback.
_MAX_FRAME_W, _MAX_FRAME_H = 100_000, 256
SYNC_PORT = 5765
HELLO_INTERVAL = 5.0 # follower broadcasts hello every 5 s
HEARTBEAT_INTERVAL = 2.0 # follower sends heartbeat every 2 s
PEER_TIMEOUT = 6.0 # leader: no heartbeat → follower gone
LEADER_TIMEOUT = 6.0 # follower: no frame → leader gone
STATUS_FILE = os.path.join(tempfile.gettempdir(), "led_matrix_sync_status.json")
class SyncRole(Enum):
STANDALONE = "standalone"
LEADER = "leader"
FOLLOWER = "follower"
class LeaderState(Enum):
NO_PEER = "no_peer"
CONNECTED = "connected"
INCOMPATIBLE = "incompatible"
class FollowerState(Enum):
STANDALONE = "standalone"
FOLLOWER = "follower"
class DisplaySyncManager:
"""
Core sync manager. Instantiated by DisplayController based on config['sync'].
The leader sends each rendered frame to the follower over UDP as raw RGB
bytes (send_frame), and for Vegas scrolling sends the whole scroll image
once per cycle as a PNG over TCP on port + 1 (send_scroll_image), then
only the scroll position. The follower draws what it receives and goes
back to its own plugins when the leader stops sending.
"""
def __init__(
self,
role_str: str,
cfg: dict,
hw_config: dict,
logger: logging.Logger,
) -> None:
"""
Args:
role_str: "standalone" | "leader" | "follower"
cfg: config['sync'] dict
hw_config: config['display']['hardware'] dict (this Pi's own config)
logger: framework logger
"""
try:
self.role = SyncRole(role_str)
except ValueError:
logger.warning("Invalid sync role '%s', defaulting to standalone", role_str)
self.role = SyncRole.STANDALONE
self.logger = logger
self.port = int(cfg.get("port", SYNC_PORT))
self._hw_config = hw_config
# Leader state
self._leader_state = LeaderState.NO_PEER
self._peer_ip: Optional[str] = None
self._peer_compatible: bool = False
self._peer_chain: int = 0
self._last_heartbeat_time: float = 0.0
self._leader_width: int = 0 # set by display_controller after init
self._oversized_frame_warned: bool = False
# Follower state
self._follower_state = FollowerState.STANDALONE
self._latest_frame: Optional[Image.Image] = None # pixel-frame fallback
self._latest_scroll_x: Optional[float] = None # Vegas scroll position
self._last_leader_frame_time: float = 0.0
self._frame_lock = threading.Lock()
self._leader_ip: Optional[str] = None
self._on_new_cycle: Optional[Callable[[], None]] = None # called when leader starts new cycle
self._on_scroll_image: Optional[Callable[[Image.Image], None]] = None # called with Image when received
self._pending_scroll_image: Optional[Image.Image] = None # image received before callback set
self._scroll_image_lock = threading.Lock() # guards _on_scroll_image / _pending_scroll_image
self._img_server_sock = None # TCP server for scroll image transfer
# Leader state additions
self._on_follower_connected: Optional[Callable[[], None]] = None # called when follower connects
self._error_message: Optional[str] = None
self._running = False
self._recv_sock: Optional[socket.socket] = None
self._send_sock: Optional[socket.socket] = None
if self.role == SyncRole.STANDALONE:
return
if self.role == SyncRole.LEADER:
self._start_leader()
elif self.role == SyncRole.FOLLOWER:
self._start_follower()
# ------------------------------------------------------------------ #
# Leader setup #
# ------------------------------------------------------------------ #
def _start_leader(self) -> None:
# Receive socket: listens for hello + heartbeat from follower
self._recv_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) # nosec B104
self._recv_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
self._recv_sock.bind(("", self.port)) # nosec B104 — intentional: must receive UDP broadcast on all interfaces
self._recv_sock.settimeout(1.0)
# Send socket: unicast frames + hello_ack to follower
self._send_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
self._running = True
threading.Thread(
target=self._leader_recv_loop, daemon=True, name="sync-leader-recv"
).start()
threading.Thread(
target=self._leader_watchdog, daemon=True, name="sync-leader-watchdog"
).start()
self.logger.info("Sync: leader started on UDP port %d", self.port)
self.write_status_file()
def _leader_recv_loop(self) -> None:
while self._running:
try:
data, addr = self._recv_sock.recvfrom(1024)
sender_ip = addr[0]
try:
msg = json.loads(data.decode("utf-8"))
except (json.JSONDecodeError, UnicodeDecodeError):
continue
t = msg.get("t")
if t == "hello":
self._handle_hello(msg, sender_ip)
elif t == "hb":
if self._peer_ip == sender_ip:
self._last_heartbeat_time = time.time()
except socket.timeout:
continue
except Exception as exc:
self.logger.debug("Sync leader recv error: %s", exc)
# Brief backoff: a socket left in a bad state raises
# immediately, which would otherwise spin this thread at
# 100% CPU logging the same error.
time.sleep(0.1)
def _handle_hello(self, msg: dict, sender_ip: str) -> None:
hw = self._hw_config
local_rows = hw.get("rows", DEFAULT_ROWS)
local_cols = hw.get("cols", DEFAULT_COLS)
peer_rows = int(msg.get("rows", 0))
peer_cols = int(msg.get("cols", 0))
peer_chain = int(msg.get("chain", DEFAULT_CHAIN_LENGTH))
compatible = peer_rows == local_rows and peer_cols == local_cols
self._peer_ip = sender_ip
self._peer_compatible = compatible
self._peer_chain = peer_chain
self._last_heartbeat_time = time.time()
prev_state = self._leader_state
if compatible:
if prev_state != LeaderState.CONNECTED:
self.logger.info(
"Sync: follower connected at %s (chain=%d)", sender_ip, peer_chain
)
self._leader_state = LeaderState.CONNECTED
self._error_message = None
# Send scroll image immediately on new connection so follower has identical content
if prev_state != LeaderState.CONNECTED and self._on_follower_connected:
threading.Thread(
target=self._on_follower_connected,
daemon=True, name="sync-leader-img-push"
).start()
else:
self._leader_state = LeaderState.INCOMPATIBLE
self._error_message = (
f"Incompatible panels: follower is {peer_cols}x{peer_rows}, "
f"leader is {local_cols}x{local_rows}. "
f"rows and cols must match between displays."
)
if prev_state != LeaderState.INCOMPATIBLE:
self.logger.error("Sync: %s", self._error_message)
if self._leader_state != prev_state:
self.write_status_file()
ack = json.dumps({
"t": "hello_ack",
"compatible": compatible,
"leader_width": self._leader_width,
"error": self._error_message,
}).encode("utf-8")
try:
self._send_sock.sendto(ack, (sender_ip, self.port))
except Exception as exc:
self.logger.debug("Sync: hello_ack send failed: %s", exc)
def _leader_watchdog(self) -> None:
while self._running:
time.sleep(1.0)
if self._leader_state == LeaderState.CONNECTED:
if time.time() - self._last_heartbeat_time > PEER_TIMEOUT:
self.logger.info(
"Sync: follower heartbeat timeout — peer disconnected"
)
self._leader_state = LeaderState.NO_PEER
self._peer_ip = None
self._peer_compatible = False
self.write_status_file()
def _image_server_loop(self) -> None:
"""Follower: TCP server that receives the leader's scroll image at each new cycle."""
while self._running:
try:
conn, addr = self._img_server_sock.accept()
conn.settimeout(10.0)
try:
# 4-byte big-endian length prefix
hdr = b""
while len(hdr) < 4:
chunk = conn.recv(4 - len(hdr))
if not chunk:
break
hdr += chunk
if len(hdr) < 4:
continue
length = int.from_bytes(hdr, "big")
_MAX_IMAGE_BYTES = 10 * 1024 * 1024 # 10 MB — well above any real scroll image
if length <= 0 or length > _MAX_IMAGE_BYTES:
self.logger.warning(
"Sync: rejected TCP image with invalid length %d (max %d) from %s",
length, _MAX_IMAGE_BYTES, addr,
)
conn.close()
continue
data = bytearray()
while len(data) < length:
chunk = conn.recv(min(65536, length - len(data)))
if not chunk:
break
data.extend(chunk)
img = Image.open(io.BytesIO(data))
if img.width > _MAX_FRAME_W or img.height > _MAX_FRAME_H:
self.logger.warning(
"Sync: rejected oversized scroll image %dx%d (max %dx%d) from %s",
img.width, img.height, _MAX_FRAME_W, _MAX_FRAME_H, addr,
)
continue
try:
img.load()
except (Image.DecompressionBombError, ValueError) as exc:
self.logger.warning("Sync: rejected decompression bomb from %s: %s", addr, exc)
continue
self.logger.info(
"Sync: received scroll image %dx%d (%d bytes compressed)",
img.width, img.height, length,
)
with self._scroll_image_lock:
if self._on_scroll_image:
cb = self._on_scroll_image
else:
# Callback not registered yet (startup race) — cache it
self._pending_scroll_image = img
cb = None
if cb:
cb(img)
finally:
conn.close()
except socket.timeout:
continue
except Exception as exc:
self.logger.debug("Sync: image server error: %s", exc)
def send_scroll_image(self, image: Image.Image) -> None:
"""Leader: send the full scroll image to the follower via TCP.
PNG compression typically reduces a 5000×32 image to ~20–50KB,
transferring in <20ms on local WiFi. Called at new_cycle and on
first connection so both Pis always have identical cached_arrays.
"""
if self.role != SyncRole.LEADER:
return
if self._leader_state != LeaderState.CONNECTED or not self._peer_ip:
return
try:
buf = io.BytesIO()
image.save(buf, format="PNG", optimize=True)
data = buf.getvalue()
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock:
sock.settimeout(5.0)
sock.connect((self._peer_ip, self.port + 1))
sock.sendall(len(data).to_bytes(4, "big") + data)
self.logger.info(
"Sync: sent scroll image %dx%d (%d bytes compressed)",
image.width, image.height, len(data),
)
except Exception as exc:
self.logger.debug("Sync: image send error: %s", exc)
def set_on_follower_connected(self, callback: Callable[[], None]) -> None:
"""Leader: callback fired (in a thread) when a compatible follower first connects.
Use this to push the current scroll image immediately.
If a follower is already connected when this is called, fires right away
(handles the race where follower connects during leader startup).
"""
self._on_follower_connected = callback
if self._leader_state == LeaderState.CONNECTED:
threading.Thread(
target=callback, daemon=True, name="sync-leader-img-push-late"
).start()
def set_on_scroll_image(self, callback: Callable[[Image.Image], None]) -> None:
"""Follower: callback fired with the received Image when leader sends scroll image.
If an image was received before this callback was registered (startup race),
fires immediately with that cached image.
"""
with self._scroll_image_lock:
self._on_scroll_image = callback
pending = self._pending_scroll_image
self._pending_scroll_image = None
if pending is not None:
callback(pending)
def send_scroll_x(self, scroll_x: float) -> None:
"""Leader (Vegas mode): broadcast scroll position instead of a pixel frame.
The follower renders from its own local pipeline at scroll_x - display_width.
~20 bytes vs ~18KB for raw frames — eliminates all content-change artifacts.
"""
if self.role != SyncRole.LEADER:
return
if self._leader_state != LeaderState.CONNECTED or not self._peer_ip:
return
try:
msg = json.dumps({"t": "sx", "x": round(scroll_x, 2)}).encode("utf-8")
self._send_sock.sendto(msg, (self._peer_ip, self.port))
except Exception as exc:
self.logger.debug("Sync: scroll_x send error: %s", exc)
def send_new_cycle(self) -> None:
"""Leader: signal that a new scroll cycle has started so follower rebuilds its image."""
if self.role != SyncRole.LEADER:
return
if self._leader_state != LeaderState.CONNECTED or not self._peer_ip:
return
try:
self._send_sock.sendto(b'{"t":"nc"}', (self._peer_ip, self.port))
except Exception as exc:
self.logger.debug("Sync: new_cycle send error: %s", exc)
def send_frame(self, image: Image.Image) -> None:
"""Leader: send a rendered frame to the follower as raw RGB bytes.
Raw format is orders of magnitude faster than PNG on Pi hardware —
no encode on sender, no decode on receiver.
Packet: 8-byte magic + 4-byte (width, height) header + raw RGB bytes.
"""
if self.role != SyncRole.LEADER:
return
if self._leader_state != LeaderState.CONNECTED or not self._peer_ip:
return
try:
arr = np.asarray(image.convert("RGB"), dtype=np.uint8)
header = _RAW_MAGIC + _RAW_HEADER.pack(image.width, image.height)
data = header + arr.tobytes()
if len(data) <= 65000:
self._send_sock.sendto(data, (self._peer_ip, self.port))
elif not self._oversized_frame_warned:
self._oversized_frame_warned = True
self.logger.warning(
"Sync: frame too large for UDP (%d bytes, max 65000) — "
"image %dx%d will not be sent; use TCP image sync instead",
len(data), image.width, image.height,
)
except Exception as exc:
self.logger.debug("Sync: frame send error: %s", exc)
def set_leader_width(self, width: int) -> None:
"""Called by DisplayController once display_manager.width is known."""
self._leader_width = width
# ------------------------------------------------------------------ #
# Follower setup #
# ------------------------------------------------------------------ #
def _start_follower(self) -> None:
# Receive socket: listens for frames + hello_ack from leader
self._recv_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
self._recv_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
self._recv_sock.bind(("", self.port)) # nosec B104 — intentional: must receive UDP broadcast on all interfaces
self._recv_sock.settimeout(0.1)
# Send socket: broadcasts hello + heartbeat
self._send_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
self._send_sock.setsockopt(socket.SOL_SOCKET, socket.SO_BROADCAST, 1)
self._running = True
threading.Thread(
target=self._follower_recv_loop, daemon=True, name="sync-follower-recv"
).start()
threading.Thread(
target=self._follower_announce_loop, daemon=True, name="sync-follower-announce"
).start()
threading.Thread(
target=self._follower_watchdog, daemon=True, name="sync-follower-watchdog"
).start()
# TCP server: receives scroll images from leader (port + 1)
self._img_server_sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) # nosec B104
self._img_server_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
self._img_server_sock.bind(("", self.port + 1)) # nosec B104 — intentional: TCP server must accept connections on all interfaces
self._img_server_sock.listen(1)
self._img_server_sock.settimeout(1.0)
threading.Thread(
target=self._image_server_loop, daemon=True, name="sync-image-server"
).start()
self.logger.info(
"Sync: follower started on UDP port %d, image server on TCP %d",
self.port, self.port + 1,
)
self.write_status_file()
def _handle_received_frame(self, img: Image.Image, sender_ip: str) -> None:
"""Record a decoded leader frame and enter follower mode if needed."""
with self._frame_lock:
self._latest_frame = img
self._enter_follower_mode(sender_ip)
def _enter_follower_mode(self, sender_ip: str) -> bool:
"""Note that the leader at ``sender_ip`` just sent something, and
switch from standalone to follower mode if not already following.
Returns True if this call made the switch."""
self._last_leader_frame_time = time.time()
self._leader_ip = sender_ip
if self._follower_state != FollowerState.STANDALONE:
return False
self._follower_state = FollowerState.FOLLOWER
self.logger.info(
"Sync: leader active at %s — switching to follower mode",
sender_ip,
)
self.write_status_file()
return True
def _follower_recv_loop(self) -> None:
while self._running:
try:
data, addr = self._recv_sock.recvfrom(65535)
sender_ip = addr[0]
if data[:8] == _RAW_MAGIC:
# Magic-tagged raw RGB frame — self-describing, no guessing.
try:
w, h = _RAW_HEADER.unpack(data[8:12])
raw = data[12:]
img = Image.frombuffer(
"RGB", (w, h), raw, "raw", "RGB", 0, 1
)
self._handle_received_frame(img, sender_ip)
except Exception as exc:
self.logger.debug("Sync: frame decode error: %s", exc)
else:
# No magic prefix. Whether the payload parses as JSON
# decides between a control message and a legacy
# (pre-magic) PNG frame — both wire formats are
# self-describing, so no size heuristic is needed. A
# >512-byte control message used to be misrouted into
# image decode and silently dropped.
try:
msg = json.loads(data.decode("utf-8"))
except (json.JSONDecodeError, UnicodeDecodeError):
# Not JSON — try a legacy PNG frame.
try:
img = Image.open(io.BytesIO(data))
if img.width > _MAX_FRAME_W or img.height > _MAX_FRAME_H:
# Same cap the TCP image path applies: decode
# is deferred until load(), so check first.
self.logger.debug(
"Sync: rejected oversized legacy frame %dx%d from %s",
img.width, img.height, sender_ip,
)
continue
img.load()
self._handle_received_frame(img, sender_ip)
except Exception as exc:
self.logger.debug("Sync: frame decode error: %s", exc)
continue
# It parsed, so it is a control message and never a
# frame. Read and validate its fields under a guard —
# a UDP payload is attacker-shaped, so a non-object
# body makes .get() raise AttributeError and an "sx"
# carrying a non-numeric x raises ValueError/TypeError
# — but dispatch the callback *outside* it. Running
# the callback in here would let a fault in someone
# else's code read as a malformed packet and be
# logged as one.
fire_new_cycle = False
try:
t = msg.get("t")
if t == "hello_ack":
self._leader_ip = sender_ip
self._peer_compatible = msg.get("compatible", False)
self._error_message = msg.get("error")
if not self._peer_compatible and self._error_message:
self.logger.error(
"Sync: leader rejected handshake — %s",
self._error_message,
)
self.write_status_file()
elif t == "sx":
# Vegas scroll-position sync — tiny message, renders locally
scroll_x = float(msg["x"])
if not math.isfinite(scroll_x):
# json.loads accepts the NaN/Infinity literals,
# and float("nan") accepts the strings, so a
# non-finite x reaches here intact. Left alone
# it poisons every offset computed from it —
# NaN comparisons are all false, so the
# follower renders a frame it can never scroll
# back from. Treat it as malformed.
raise ValueError(f"non-finite scroll x: {msg['x']!r}")
self._latest_scroll_x = scroll_x
if self._enter_follower_mode(sender_ip):
fire_new_cycle = True # build initial scroll image
elif t == "nc":
# Leader started a new scroll cycle — rebuild local image
fire_new_cycle = True
except (KeyError, AttributeError, TypeError, ValueError) as exc:
self.logger.debug("Sync: malformed control message: %s", exc)
continue
if fire_new_cycle and self._on_new_cycle:
self._on_new_cycle()
except socket.timeout:
continue
except Exception as exc:
self.logger.debug("Sync follower recv error: %s", exc)
time.sleep(0.1)
def _follower_announce_loop(self) -> None:
hw = self._hw_config
hello = json.dumps({
"t": "hello",
"rows": hw.get("rows", DEFAULT_ROWS),
"cols": hw.get("cols", DEFAULT_COLS),
"chain": hw.get("chain_length", DEFAULT_CHAIN_LENGTH),
}).encode("utf-8")
heartbeat = json.dumps({"t": "hb"}).encode("utf-8")
dest = ("<broadcast>", self.port)
last_hello = 0.0
last_hb = 0.0
while self._running:
now = time.time()
if now - last_hello >= HELLO_INTERVAL:
try:
self._send_sock.sendto(hello, dest)
last_hello = now
except Exception as exc:
self.logger.debug("Sync: hello broadcast error: %s", exc)
if now - last_hb >= HEARTBEAT_INTERVAL:
try:
self._send_sock.sendto(heartbeat, dest)
last_hb = now
except Exception as exc:
self.logger.debug("Sync: heartbeat error: %s", exc)
time.sleep(0.5)
def _follower_watchdog(self) -> None:
while self._running:
time.sleep(1.0)
if self._follower_state == FollowerState.FOLLOWER:
if time.time() - self._last_leader_frame_time > LEADER_TIMEOUT:
self.logger.info(
"Sync: leader frame timeout — returning to standalone mode"
)
self._follower_state = FollowerState.STANDALONE
with self._frame_lock:
self._latest_frame = None
self.write_status_file()
# ------------------------------------------------------------------ #
# Public API #
# ------------------------------------------------------------------ #
def is_follower_active(self) -> bool:
"""True when this Pi is in active follower mode (receiving frames)."""
return (
self.role == SyncRole.FOLLOWER
and self._follower_state == FollowerState.FOLLOWER
)
def get_latest_scroll_x(self) -> Optional[float]:
"""Follower: return the most recently received Vegas scroll position, or None."""
return self._latest_scroll_x
def set_on_new_cycle(self, callback: Callable[[], None]) -> None:
"""Follower: register a callback fired when the leader starts a new scroll cycle.
Used to trigger a local start_new_cycle() so both Pis rebuild from same fresh data.
"""
self._on_new_cycle = callback
def get_latest_frame(self) -> Optional[Image.Image]:
"""Follower: return the most recently received pixel frame (non-Vegas fallback)."""
with self._frame_lock:
return self._latest_frame
def get_status(self) -> dict:
"""Return sync state dict for the web API status endpoint."""
hw = self._hw_config
base = {
"role": self.role.value,
"port": self.port,
"local_rows": hw.get("rows", DEFAULT_ROWS),
"local_cols": hw.get("cols", DEFAULT_COLS),
"local_chain": hw.get("chain_length", DEFAULT_CHAIN_LENGTH),
}
if self.role == SyncRole.STANDALONE:
return {**base, "state": "standalone"}
if self.role == SyncRole.LEADER:
return {
**base,
"state": self._leader_state.value,
"peer_ip": self._peer_ip,
"peer_compatible": self._peer_compatible,
"peer_chain": self._peer_chain,
"leader_width": self._leader_width,
"error": self._error_message,
}
# Follower
return {
**base,
"state": self._follower_state.value,
"leader_ip": self._leader_ip,
"peer_compatible": self._peer_compatible,
"error": self._error_message,
}
def write_status_file(self) -> None:
"""Write current sync status to STATUS_FILE for the web UI to read."""
try:
status = self.get_status()
status["ts"] = time.time()
tmp = STATUS_FILE + ".tmp"
with open(tmp, "w") as f:
json.dump(status, f)
os.replace(tmp, STATUS_FILE)
except Exception as exc:
self.logger.debug("Sync: status file write error: %s", exc)
def stop(self) -> None:
"""Shut down threads and close sockets."""
self._running = False
for sock in (self._recv_sock, self._send_sock, self._img_server_sock):
if sock:
try:
sock.close()
except Exception as exc:
self.logger.debug("Sync: error closing socket: %s", exc)