Files
LEDMatrix/test/test_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

1097 lines
44 KiB
Python

"""
Tests for src/common/sync_manager.py — the UDP leader/follower protocol
that synchronizes scrolling content across two LED matrix displays.
This module had zero coverage: it only ever appeared in the suite as a
MagicMock() stand-in (test_vegas_continuous_refresh.py,
test_display_controller_vegas_tick.py), so none of its real framing,
handshake, or socket logic was exercised.
Most tests build the manager via object.__new__() + manual attribute
assignment (the test_display_controller_vegas_tick.py bare-stub pattern)
so no real sockets open and no background threads start. Receive loops are
driven synchronously by once_then_stop(): the mocked socket call returns
one crafted packet, then flips _running False and raises socket.timeout,
so `while self._running:` exits after exactly one real iteration.
Regression coverage for three fixed bugs:
- Both recv loops' generic `except Exception` retried with no delay, so a
socket stuck raising a non-timeout error spun the thread at 100% CPU.
- _follower_recv_loop dispatched on `data[:8] == _RAW_MAGIC or
len(data) > 512`, which routed any control message over 512 bytes into
the image decoder (dropping it) and any raw frame under 512 bytes into
the JSON parser.
- _oversized_frame_warned was read via getattr(self, ..., False) instead of
being initialized in __init__.
"""
import io
import json
import socket
import threading
import time
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import MagicMock, patch
import numpy as np
import pytest
from PIL import Image
from test._api_v3_test_helpers import api_v3_client, api_v3_module # noqa: F401
from src.common import sync_manager
from src.common.sync_manager import (
DisplaySyncManager,
FollowerState,
LeaderState,
SyncRole,
)
@pytest.fixture(autouse=True)
def _isolated_status_file(tmp_path, monkeypatch):
# STATUS_FILE is a module-level fixed path under tempfile.gettempdir() —
# genuinely shared state between tests and even between processes.
monkeypatch.setattr(
sync_manager, "STATUS_FILE", str(tmp_path / "led_matrix_sync_status.json"))
def make_manager(role=SyncRole.STANDALONE, hw_config=None):
"""Bare stub bypassing __init__'s socket/thread setup."""
mgr = object.__new__(DisplaySyncManager)
mgr.role = role
mgr.logger = MagicMock()
mgr.port = sync_manager.SYNC_PORT
mgr._hw_config = hw_config or {"rows": 32, "cols": 64, "chain_length": 1}
mgr._leader_state = LeaderState.NO_PEER
mgr._peer_ip = None
mgr._peer_compatible = False
mgr._peer_chain = 0
mgr._last_heartbeat_time = 0.0
mgr._leader_width = 0
mgr._oversized_frame_warned = False
mgr._follower_state = FollowerState.STANDALONE
mgr._latest_frame = None
mgr._latest_scroll_x = None
mgr._last_leader_frame_time = 0.0
mgr._frame_lock = threading.Lock()
mgr._leader_ip = None
mgr._on_new_cycle = None
mgr._on_scroll_image = None
mgr._pending_scroll_image = None
mgr._scroll_image_lock = threading.Lock()
mgr._img_server_sock = None
mgr._on_follower_connected = None
mgr._error_message = None
mgr._running = False
mgr._recv_sock = None
mgr._send_sock = None
return mgr
def once_then_stop(mgr, value):
"""side_effect returning `value` once, then stopping the enclosing loop."""
state = {"served": False}
def _side_effect(*args, **kwargs):
if not state["served"]:
state["served"] = True
return value
mgr._running = False
raise socket.timeout()
return _side_effect
def raise_n_then_stop(mgr, exc, count):
"""side_effect raising `exc` `count` times, then stopping the loop."""
state = {"n": 0}
def _side_effect(*args, **kwargs):
state["n"] += 1
if state["n"] <= count:
raise exc
mgr._running = False
raise socket.timeout()
return _side_effect
def fake_clock(monkeypatch, *, time_fn=None, sleep_fn=None):
"""Swap sync_manager's own `time` reference for a private stand-in.
sync_manager.time IS the stdlib module, so patching attributes on it
would freeze the clock and no-op sleep for the whole process —
including the daemon threads earlier tests left running, which is a
hard-to-trace source of cross-test flakiness. Rebinding the module's
reference keeps the patch scoped to the code under test. Anything not
overridden falls through to the real functions.
"""
monkeypatch.setattr(sync_manager, "time", SimpleNamespace(
time=time_fn or time.time,
sleep=sleep_fn or time.sleep,
))
def run_watchdog_once(monkeypatch, mgr, watchdog, now):
"""Run exactly one watchdog iteration at a frozen wall-clock time."""
fake_clock(monkeypatch,
time_fn=lambda: now,
sleep_fn=lambda _: setattr(mgr, "_running", False))
mgr._running = True
watchdog()
class FakeConn:
"""Minimal TCP connection stand-in whose recv() drains a byte buffer."""
def __init__(self, payload: bytes):
self._buf = payload
self.closed = False
def settimeout(self, _):
pass
def recv(self, n):
chunk, self._buf = self._buf[:n], self._buf[n:]
return chunk
def close(self):
self.closed = True
def png_bytes(size=(10, 10), color=(1, 2, 3)) -> bytes:
buf = io.BytesIO()
Image.new("RGB", size, color).save(buf, format="PNG")
return buf.getvalue()
def raw_frame_packet(width, height, color=(10, 20, 30)) -> bytes:
arr = np.asarray(Image.new("RGB", (width, height), color), dtype=np.uint8)
return _magic_header(width, height) + arr.tobytes()
def _magic_header(width, height) -> bytes:
return sync_manager._RAW_MAGIC + sync_manager._RAW_HEADER.pack(width, height)
def length_prefixed(payload: bytes) -> bytes:
return len(payload).to_bytes(4, "big") + payload
class TestRoleParsing:
def test_leader_role(self, monkeypatch):
monkeypatch.setattr(DisplaySyncManager, "_start_leader", lambda self: None)
assert DisplaySyncManager("leader", {}, {}, MagicMock()).role is SyncRole.LEADER
def test_follower_role(self, monkeypatch):
monkeypatch.setattr(DisplaySyncManager, "_start_follower", lambda self: None)
assert DisplaySyncManager("follower", {}, {}, MagicMock()).role is SyncRole.FOLLOWER
def test_standalone_starts_nothing(self):
mgr = DisplaySyncManager("standalone", {}, {}, MagicMock())
assert mgr.role is SyncRole.STANDALONE
assert mgr._running is False
assert mgr._recv_sock is None
def test_invalid_role_warns_and_falls_back(self):
logger = MagicMock()
assert DisplaySyncManager("bogus", {}, {}, logger).role is SyncRole.STANDALONE
assert logger.warning.called
def test_role_matching_is_case_sensitive(self):
# Pinned: SyncRole's values are lowercase, so "LEADER" is not
# normalized — it is simply invalid and falls back to standalone.
logger = MagicMock()
assert DisplaySyncManager("LEADER", {}, {}, logger).role is SyncRole.STANDALONE
assert logger.warning.called
def test_port_defaults_to_module_constant(self):
assert DisplaySyncManager("standalone", {}, {}, MagicMock()).port == sync_manager.SYNC_PORT
def test_port_read_from_config(self):
assert DisplaySyncManager("standalone", {"port": 9999}, {}, MagicMock()).port == 9999
def test_oversized_frame_warned_initialized_in_init(self, monkeypatch):
# Regression: this attribute was only ever created on first use via
# getattr(self, '_oversized_frame_warned', False).
monkeypatch.setattr(DisplaySyncManager, "_start_leader", lambda self: None)
mgr = DisplaySyncManager("leader", {}, {}, MagicMock())
assert mgr._oversized_frame_warned is False
class TestHandleHello:
def test_matching_panels_connect(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._send_sock = MagicMock()
mgr._handle_hello({"t": "hello", "rows": 32, "cols": 64, "chain": 3}, "10.0.0.5")
assert mgr._leader_state is LeaderState.CONNECTED
assert mgr._peer_ip == "10.0.0.5"
assert mgr._peer_compatible is True
assert mgr._peer_chain == 3
assert mgr._error_message is None
def test_ack_reports_compatibility(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._send_sock = MagicMock()
mgr._leader_width = 128
mgr._handle_hello({"t": "hello", "rows": 32, "cols": 64, "chain": 1}, "10.0.0.5")
payload, dest = mgr._send_sock.sendto.call_args[0]
ack = json.loads(payload.decode("utf-8"))
assert ack["compatible"] is True
assert ack["leader_width"] == 128
assert dest == ("10.0.0.5", mgr.port)
def test_mismatched_panels_are_incompatible(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._send_sock = MagicMock()
mgr._handle_hello({"t": "hello", "rows": 16, "cols": 32, "chain": 1}, "10.0.0.5")
assert mgr._leader_state is LeaderState.INCOMPATIBLE
assert "Incompatible panels" in mgr._error_message
ack = json.loads(mgr._send_sock.sendto.call_args[0][0].decode("utf-8"))
assert ack["compatible"] is False
assert ack["error"] == mgr._error_message
def test_chain_length_may_differ(self):
# Documented rule: rows/cols must match, chain_length need not.
mgr = make_manager(role=SyncRole.LEADER, hw_config={"rows": 32, "cols": 64, "chain_length": 1})
mgr._send_sock = MagicMock()
mgr._handle_hello({"t": "hello", "rows": 32, "cols": 64, "chain": 4}, "10.0.0.5")
assert mgr._leader_state is LeaderState.CONNECTED
def test_connect_callback_fires_only_on_first_transition(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._send_sock = MagicMock()
fired = threading.Event()
calls = []
mgr._on_follower_connected = lambda: (calls.append(1), fired.set())
hello = {"t": "hello", "rows": 32, "cols": 64, "chain": 1}
mgr._handle_hello(hello, "10.0.0.5")
assert fired.wait(timeout=1)
assert len(calls) == 1
fired.clear()
mgr._handle_hello(hello, "10.0.0.5") # already CONNECTED
assert not fired.wait(timeout=0.2)
assert len(calls) == 1
def test_ack_send_failure_is_swallowed(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._send_sock = MagicMock()
mgr._send_sock.sendto.side_effect = OSError("network unreachable")
mgr._handle_hello({"t": "hello", "rows": 32, "cols": 64, "chain": 1}, "10.0.0.5")
assert mgr._leader_state is LeaderState.CONNECTED # state still updated
assert mgr.logger.debug.called
class TestWatchdogs:
def test_leader_drops_peer_after_heartbeat_timeout(self, monkeypatch):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
mgr._peer_ip = "10.0.0.1"
mgr._peer_compatible = True
mgr._last_heartbeat_time = 0.0
run_watchdog_once(monkeypatch, mgr, mgr._leader_watchdog,
now=sync_manager.PEER_TIMEOUT + 1)
assert mgr._leader_state is LeaderState.NO_PEER
assert mgr._peer_ip is None
assert mgr._peer_compatible is False
def test_leader_keeps_peer_within_timeout(self, monkeypatch):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
mgr._peer_ip = "10.0.0.1"
mgr._last_heartbeat_time = 100.0
run_watchdog_once(monkeypatch, mgr, mgr._leader_watchdog, now=101.0)
assert mgr._leader_state is LeaderState.CONNECTED
assert mgr._peer_ip == "10.0.0.1"
def test_leader_watchdog_ignores_disconnected_state(self, monkeypatch):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.INCOMPATIBLE
mgr._last_heartbeat_time = 0.0
run_watchdog_once(monkeypatch, mgr, mgr._leader_watchdog, now=10_000)
assert mgr._leader_state is LeaderState.INCOMPATIBLE
def test_follower_returns_to_standalone_after_frame_timeout(self, monkeypatch):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._follower_state = FollowerState.FOLLOWER
mgr._last_leader_frame_time = 0.0
mgr._latest_frame = Image.new("RGB", (2, 2))
run_watchdog_once(monkeypatch, mgr, mgr._follower_watchdog,
now=sync_manager.LEADER_TIMEOUT + 1)
assert mgr._follower_state is FollowerState.STANDALONE
assert mgr.get_latest_frame() is None
def test_follower_keeps_frames_within_timeout(self, monkeypatch):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._follower_state = FollowerState.FOLLOWER
mgr._last_leader_frame_time = 100.0
mgr._latest_frame = Image.new("RGB", (2, 2))
run_watchdog_once(monkeypatch, mgr, mgr._follower_watchdog, now=101.0)
assert mgr._follower_state is FollowerState.FOLLOWER
assert mgr.get_latest_frame() is not None
class TestLeaderRecvLoop:
def _drive(self, mgr, payload, sender="10.0.0.8"):
mgr._recv_sock = MagicMock()
mgr._recv_sock.recvfrom.side_effect = once_then_stop(mgr, (payload, (sender, 1)))
mgr._running = True
mgr._leader_recv_loop()
def test_hello_is_dispatched(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._send_sock = MagicMock()
self._drive(mgr, json.dumps(
{"t": "hello", "rows": 32, "cols": 64, "chain": 1}).encode())
assert mgr._leader_state is LeaderState.CONNECTED
assert mgr._peer_ip == "10.0.0.8"
def test_heartbeat_from_known_peer_refreshes_timer(self, monkeypatch):
mgr = make_manager(role=SyncRole.LEADER)
mgr._peer_ip = "10.0.0.8"
fake_clock(monkeypatch, time_fn=lambda: 12345.0)
self._drive(mgr, json.dumps({"t": "hb"}).encode())
assert mgr._last_heartbeat_time == 12345.0
def test_heartbeat_from_stranger_is_ignored(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._peer_ip = "10.0.0.8"
mgr._last_heartbeat_time = 5.0
self._drive(mgr, json.dumps({"t": "hb"}).encode(), sender="10.0.0.99")
assert mgr._last_heartbeat_time == 5.0
def test_unknown_message_type_ignored(self):
mgr = make_manager(role=SyncRole.LEADER)
self._drive(mgr, json.dumps({"t": "who-knows"}).encode())
assert mgr._leader_state is LeaderState.NO_PEER
def test_malformed_json_is_swallowed(self):
mgr = make_manager(role=SyncRole.LEADER)
self._drive(mgr, b"{not json")
assert mgr._leader_state is LeaderState.NO_PEER
def test_undecodable_bytes_are_swallowed(self):
mgr = make_manager(role=SyncRole.LEADER)
self._drive(mgr, b"\xff\xfe\x00bad")
assert mgr._leader_state is LeaderState.NO_PEER
def test_backs_off_between_repeated_errors(self, monkeypatch):
# Regression: without a sleep this loop spun at 100% CPU whenever
# the socket raised a non-timeout error on every call.
mgr = make_manager(role=SyncRole.LEADER)
mgr._recv_sock = MagicMock()
mgr._recv_sock.recvfrom.side_effect = raise_n_then_stop(mgr, OSError("boom"), 3)
sleeps = MagicMock()
fake_clock(monkeypatch, sleep_fn=sleeps)
mgr._running = True
mgr._leader_recv_loop()
assert sleeps.call_count == 3
sleeps.assert_called_with(0.1)
class TestFollowerRecvLoop:
def _drive(self, mgr, payload, sender="10.0.0.2"):
mgr._recv_sock = MagicMock()
mgr._recv_sock.recvfrom.side_effect = once_then_stop(mgr, (payload, (sender, 1)))
mgr._running = True
mgr._follower_recv_loop()
def test_small_raw_frame_is_decoded(self):
# Regression: a raw frame under the old 512-byte threshold was sent
# to the JSON parser and dropped.
mgr = make_manager(role=SyncRole.FOLLOWER)
packet = raw_frame_packet(4, 3)
assert len(packet) <= 512
self._drive(mgr, packet)
frame = mgr.get_latest_frame()
assert frame is not None and frame.size == (4, 3)
assert mgr._follower_state is FollowerState.FOLLOWER
def test_large_raw_frame_is_decoded(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
packet = raw_frame_packet(64, 32)
assert len(packet) > 512
self._drive(mgr, packet)
assert mgr.get_latest_frame().size == (64, 32)
def test_large_control_message_is_not_routed_to_image_decode(self):
# Regression: the old `len(data) > 512` branch treated any large
# control message as frame data and silently discarded it.
mgr = make_manager(role=SyncRole.FOLLOWER)
long_error = "x" * 600
payload = json.dumps(
{"t": "hello_ack", "compatible": False, "error": long_error}).encode()
assert len(payload) > 512
self._drive(mgr, payload, sender="10.0.0.9")
assert mgr._leader_ip == "10.0.0.9"
assert mgr._peer_compatible is False
assert mgr._error_message == long_error
assert mgr.get_latest_frame() is None
assert mgr.logger.error.called
def test_legacy_png_frame_without_magic_is_decoded(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
self._drive(mgr, png_bytes(size=(5, 5)))
frame = mgr.get_latest_frame()
assert frame is not None and frame.size == (5, 5)
assert mgr._follower_state is FollowerState.FOLLOWER
def test_truncated_raw_frame_is_swallowed(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
self._drive(mgr, _magic_header(64, 32) + b"\x00" * 10) # far too short
assert mgr.get_latest_frame() is None
assert mgr.logger.debug.called
def test_garbage_payload_is_swallowed(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
self._drive(mgr, b"neither json nor a png, just bytes 1234567890")
assert mgr.get_latest_frame() is None
def test_hello_ack_updates_peer_state(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
self._drive(mgr, json.dumps(
{"t": "hello_ack", "compatible": True, "error": None}).encode(),
sender="10.0.0.6")
assert mgr._leader_ip == "10.0.0.6"
assert mgr._peer_compatible is True
assert mgr.logger.error.called is False
def test_scroll_x_switches_to_follower_and_builds_cycle(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
calls = []
mgr._on_new_cycle = lambda: calls.append(1)
self._drive(mgr, json.dumps({"t": "sx", "x": 12.34}).encode())
assert mgr._follower_state is FollowerState.FOLLOWER
assert mgr.get_latest_scroll_x() == 12.34
assert calls == [1]
def test_scroll_x_while_already_following_does_not_rebuild(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._follower_state = FollowerState.FOLLOWER
calls = []
mgr._on_new_cycle = lambda: calls.append(1)
self._drive(mgr, json.dumps({"t": "sx", "x": 5.0}).encode())
assert mgr.get_latest_scroll_x() == 5.0
assert calls == []
def test_new_cycle_message_triggers_callback(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._follower_state = FollowerState.FOLLOWER
calls = []
mgr._on_new_cycle = lambda: calls.append(1)
self._drive(mgr, json.dumps({"t": "nc"}).encode())
assert calls == [1]
def test_non_object_json_does_not_reach_the_outer_handler(self):
# A bare JSON scalar parses, then msg.get() raises AttributeError.
# That has to be caught here so the payload still gets its shot at
# the legacy-PNG fallback; escaping to the outer handler would also
# charge one malformed packet the 0.1s error backoff.
mgr = make_manager(role=SyncRole.FOLLOWER)
sleeps = MagicMock()
with patch.object(sync_manager, "time",
SimpleNamespace(time=time.time, sleep=sleeps)):
self._drive(mgr, b"12345")
assert mgr.get_latest_frame() is None
sleeps.assert_not_called()
def test_non_numeric_scroll_x_does_not_reach_the_outer_handler(self):
# float("a") raises ValueError; {"x": null} raises TypeError.
for payload in ({"t": "sx", "x": "a"}, {"t": "sx", "x": None}):
mgr = make_manager(role=SyncRole.FOLLOWER)
sleeps = MagicMock()
with patch.object(sync_manager, "time",
SimpleNamespace(time=time.time, sleep=sleeps)):
self._drive(mgr, json.dumps(payload).encode())
assert mgr.get_latest_scroll_x() is None
sleeps.assert_not_called()
def test_callback_failure_is_not_mistaken_for_a_malformed_packet(self, monkeypatch):
# A payload that parses is a control message, full stop. If the
# callback it triggers raises one of the types the field guard
# catches, that fault belongs to the callback: it must not send
# the packet to the image decoder, which would report it as a
# decode error and bury the real cause. The loop still survives
# it — the outer handler catches it like any other fault.
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._follower_state = FollowerState.FOLLOWER
def boom():
raise ValueError("callback is broken")
mgr._on_new_cycle = boom
fake_clock(monkeypatch, sleep_fn=MagicMock())
self._drive(mgr, json.dumps({"t": "nc"}).encode())
logged = " | ".join(str(c) for c in mgr.logger.debug.call_args_list)
assert "callback is broken" in logged
assert "frame decode error" not in logged
assert "malformed control message" not in logged
def test_oversized_legacy_frame_is_rejected_before_decode(self, monkeypatch):
# The UDP path is reachable by any host on the LAN, so it caps
# dimensions before load() just as the TCP image server does.
mgr = make_manager(role=SyncRole.FOLLOWER)
class Huge:
width, height = 10, sync_manager._MAX_FRAME_H + 1
def load(self):
raise AssertionError("load() must not run past the cap")
# Rebind the module's reference rather than mutating PIL.Image
# itself, which would hand Huge() to every caller in the process
# — including daemon threads earlier tests left running. Same
# reasoning as fake_clock above. The other names the receive loop
# reads off this reference pass through to the real module.
monkeypatch.setattr(sync_manager, "Image", SimpleNamespace(
open=lambda *a, **kw: Huge(),
frombuffer=Image.frombuffer,
DecompressionBombError=Image.DecompressionBombError,
))
self._drive(mgr, b"\x89PNG not really but not JSON either")
assert mgr.get_latest_frame() is None
@pytest.mark.parametrize("literal", ["NaN", "Infinity", "-Infinity"])
def test_non_finite_scroll_x_is_rejected(self, literal):
# json.loads accepts these bare literals, and float() accepts them
# as strings, so they arrive as real floats rather than raising.
# NaN in particular survives every comparison the scroll code makes
# (all false), so the follower would sit on a position it can never
# advance past. It has to be treated as a malformed message.
for payload in (b'{"t": "sx", "x": ' + literal.encode() + b'}',
json.dumps({"t": "sx", "x": literal}).encode()):
mgr = make_manager(role=SyncRole.FOLLOWER)
calls = []
mgr._on_new_cycle = lambda: calls.append(1)
self._drive(mgr, payload)
assert mgr.get_latest_scroll_x() is None
assert mgr._follower_state is FollowerState.STANDALONE
assert calls == []
def test_non_finite_scroll_x_leaves_a_good_value_in_place(self):
# The reject must not clear the last usable position either — a
# follower mid-scroll keeps rendering from where it was.
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._follower_state = FollowerState.FOLLOWER
self._drive(mgr, json.dumps({"t": "sx", "x": 7.5}).encode())
assert mgr.get_latest_scroll_x() == 7.5
self._drive(mgr, b'{"t": "sx", "x": NaN}')
assert mgr.get_latest_scroll_x() == 7.5
def test_scroll_x_missing_key_is_swallowed(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
self._drive(mgr, json.dumps({"t": "sx"}).encode()) # no "x"
assert mgr.get_latest_scroll_x() is None
def test_backs_off_between_repeated_errors(self, monkeypatch):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._recv_sock = MagicMock()
mgr._recv_sock.recvfrom.side_effect = raise_n_then_stop(mgr, OSError("boom"), 3)
sleeps = MagicMock()
fake_clock(monkeypatch, sleep_fn=sleeps)
mgr._running = True
mgr._follower_recv_loop()
assert sleeps.call_count == 3
sleeps.assert_called_with(0.1)
class TestSendFrame:
def _connected_leader(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
mgr._peer_ip = "10.0.0.1"
mgr._send_sock = MagicMock()
return mgr
def test_frame_sent_with_magic_header(self):
mgr = self._connected_leader()
mgr.send_frame(Image.new("RGB", (8, 8)))
packet = mgr._send_sock.sendto.call_args[0][0]
assert packet[:8] == sync_manager._RAW_MAGIC
assert sync_manager._RAW_HEADER.unpack(packet[8:12]) == (8, 8)
def test_oversized_frame_warns_once_and_is_dropped(self):
mgr = self._connected_leader()
big = Image.new("RGB", (300, 300)) # 270000 bytes > 65000 UDP cap
mgr.send_frame(big)
assert mgr._oversized_frame_warned is True
assert mgr.logger.warning.call_count == 1
assert not mgr._send_sock.sendto.called
mgr.send_frame(big)
assert mgr.logger.warning.call_count == 1 # still warned only once
def test_not_sent_when_no_peer(self):
mgr = self._connected_leader()
mgr._leader_state = LeaderState.NO_PEER
mgr.send_frame(Image.new("RGB", (8, 8)))
assert not mgr._send_sock.sendto.called
def test_follower_never_sends(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._send_sock = MagicMock()
mgr.send_frame(Image.new("RGB", (8, 8)))
assert not mgr._send_sock.sendto.called
def test_send_error_is_swallowed(self):
mgr = self._connected_leader()
mgr._send_sock.sendto.side_effect = OSError("no route")
mgr.send_frame(Image.new("RGB", (8, 8))) # must not raise
assert mgr.logger.debug.called
class TestSendControlMessages:
def _connected_leader(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
mgr._peer_ip = "10.0.0.1"
mgr._send_sock = MagicMock()
return mgr
def test_send_scroll_x_rounds_to_two_places(self):
mgr = self._connected_leader()
mgr.send_scroll_x(3.14159)
msg = json.loads(mgr._send_sock.sendto.call_args[0][0].decode())
assert msg == {"t": "sx", "x": 3.14}
def test_send_new_cycle(self):
mgr = self._connected_leader()
mgr.send_new_cycle()
msg = json.loads(mgr._send_sock.sendto.call_args[0][0].decode())
assert msg == {"t": "nc"}
def test_control_messages_noop_when_disconnected(self):
mgr = self._connected_leader()
mgr._leader_state = LeaderState.NO_PEER
mgr.send_scroll_x(1.0)
mgr.send_new_cycle()
assert not mgr._send_sock.sendto.called
def test_set_leader_width(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr.set_leader_width(256)
assert mgr._leader_width == 256
class TestImageServerLoop:
def _drive(self, mgr, conn):
mgr._img_server_sock = MagicMock()
mgr._img_server_sock.accept.side_effect = once_then_stop(
mgr, (conn, ("10.0.0.1", 1)))
mgr._running = True
mgr._image_server_loop()
def test_rejects_non_positive_length(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._on_scroll_image = MagicMock()
self._drive(mgr, FakeConn((0).to_bytes(4, "big")))
assert mgr.logger.warning.called
mgr._on_scroll_image.assert_not_called()
def test_rejects_oversized_length(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._on_scroll_image = MagicMock()
self._drive(mgr, FakeConn((11 * 1024 * 1024).to_bytes(4, "big")))
assert mgr.logger.warning.called
mgr._on_scroll_image.assert_not_called()
def test_rejects_oversized_dimensions(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._on_scroll_image = MagicMock()
self._drive(mgr, FakeConn(length_prefixed(png_bytes(size=(300, 300)))))
assert mgr.logger.warning.called
mgr._on_scroll_image.assert_not_called()
def test_rejects_decompression_bomb(self, monkeypatch):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._on_scroll_image = MagicMock()
class BombImage:
width = height = 10
def load(self):
raise Image.DecompressionBombError("too many pixels")
monkeypatch.setattr(sync_manager.Image, "open", lambda *a, **kw: BombImage())
self._drive(mgr, FakeConn(length_prefixed(png_bytes())))
assert mgr.logger.warning.called
mgr._on_scroll_image.assert_not_called()
def test_valid_image_invokes_callback(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
received = []
mgr._on_scroll_image = received.append
self._drive(mgr, FakeConn(length_prefixed(png_bytes(size=(10, 10)))))
assert len(received) == 1
assert received[0].size == (10, 10)
def test_image_cached_when_callback_not_yet_registered(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._on_scroll_image = None
self._drive(mgr, FakeConn(length_prefixed(png_bytes(size=(6, 6)))))
assert mgr._pending_scroll_image is not None
assert mgr._pending_scroll_image.size == (6, 6)
def test_short_header_is_skipped(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._on_scroll_image = MagicMock()
self._drive(mgr, FakeConn(b"\x00\x01")) # under the 4-byte prefix
mgr._on_scroll_image.assert_not_called()
def test_connection_always_closed(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
conn = FakeConn(length_prefixed(png_bytes()))
self._drive(mgr, conn)
assert conn.closed is True
class TestScrollImageCallback:
def test_pending_image_delivered_on_late_registration(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
img = Image.new("RGB", (3, 3))
mgr._pending_scroll_image = img
received = []
mgr.set_on_scroll_image(received.append)
assert received == [img]
assert mgr._pending_scroll_image is None
def test_no_pending_image_means_no_immediate_call(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
received = []
mgr.set_on_scroll_image(received.append)
assert received == []
class TestFollowerConnectedCallback:
def test_fires_immediately_when_already_connected(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
fired = threading.Event()
mgr.set_on_follower_connected(fired.set)
assert fired.wait(timeout=1)
def test_does_not_fire_when_no_peer(self):
mgr = make_manager(role=SyncRole.LEADER)
fired = threading.Event()
mgr.set_on_follower_connected(fired.set)
assert not fired.wait(timeout=0.2)
class TestSendScrollImage:
def test_noop_when_not_connected(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.NO_PEER
with patch.object(sync_manager.socket, "socket") as sock:
mgr.send_scroll_image(Image.new("RGB", (4, 4)))
sock.assert_not_called()
def test_noop_for_follower_role(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
with patch.object(sync_manager.socket, "socket") as sock:
mgr.send_scroll_image(Image.new("RGB", (4, 4)))
sock.assert_not_called()
def test_sends_length_prefixed_png(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
mgr._peer_ip = "10.0.0.1"
fake_sock = MagicMock()
fake_sock.__enter__ = lambda s: s
fake_sock.__exit__ = lambda s, *a: False
with patch.object(sync_manager.socket, "socket", return_value=fake_sock):
mgr.send_scroll_image(Image.new("RGB", (4, 4)))
payload = fake_sock.sendall.call_args[0][0]
assert int.from_bytes(payload[:4], "big") == len(payload) - 4
assert payload[4:8] == b"\x89PNG"
def test_connection_error_is_swallowed(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
mgr._peer_ip = "10.0.0.1"
with patch.object(sync_manager.socket, "socket", side_effect=OSError("refused")):
mgr.send_scroll_image(Image.new("RGB", (4, 4))) # must not raise
assert mgr.logger.debug.called
class TestGetStatus:
def test_standalone_shape(self):
status = make_manager(role=SyncRole.STANDALONE).get_status()
assert status["role"] == "standalone"
assert status["state"] == "standalone"
assert status["local_rows"] == 32 and status["local_cols"] == 64
def test_leader_shape(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._leader_state = LeaderState.CONNECTED
mgr._peer_ip = "10.0.0.1"
mgr._peer_compatible = True
mgr._peer_chain = 2
mgr._leader_width = 128
status = mgr.get_status()
assert status["role"] == "leader"
assert status["state"] == "connected"
assert status["peer_ip"] == "10.0.0.1"
assert status["peer_chain"] == 2
assert status["leader_width"] == 128
def test_follower_shape(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._follower_state = FollowerState.FOLLOWER
mgr._leader_ip = "10.0.0.2"
status = mgr.get_status()
assert status["role"] == "follower"
assert status["state"] == "follower"
assert status["leader_ip"] == "10.0.0.2"
assert "peer_chain" not in status
def test_is_follower_active(self):
mgr = make_manager(role=SyncRole.FOLLOWER)
assert mgr.is_follower_active() is False
mgr._follower_state = FollowerState.FOLLOWER
assert mgr.is_follower_active() is True
def test_leader_is_never_follower_active(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._follower_state = FollowerState.FOLLOWER
assert mgr.is_follower_active() is False
class TestWriteStatusFile:
def test_writes_status_and_cleans_up_temp(self):
mgr = make_manager(role=SyncRole.STANDALONE)
mgr.write_status_file()
data = json.loads(Path(sync_manager.STATUS_FILE).read_text())
assert data["role"] == "standalone"
assert "ts" in data
assert not Path(sync_manager.STATUS_FILE + ".tmp").exists()
def test_write_failure_is_swallowed(self, monkeypatch):
mgr = make_manager(role=SyncRole.STANDALONE)
monkeypatch.setattr("builtins.open", MagicMock(side_effect=OSError("disk full")))
mgr.write_status_file() # must not raise
assert mgr.logger.debug.called
def test_web_status_endpoint_reads_the_file_that_was_written(self, api_v3_client):
"""GET /sync/status reads STATUS_FILE, which lives under
tempfile.gettempdir() -- not always /tmp."""
mgr = make_manager(role=SyncRole.LEADER)
mgr.write_status_file()
response = api_v3_client.get("/api/v3/sync/status")
assert response.get_json()["data"]["role"] == "leader"
class TestStop:
def _stub_with_sockets(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._recv_sock = MagicMock()
mgr._send_sock = MagicMock()
mgr._img_server_sock = MagicMock()
return mgr
def test_closes_every_socket(self):
mgr = self._stub_with_sockets()
mgr.stop()
assert mgr._running is False
mgr._recv_sock.close.assert_called_once()
mgr._send_sock.close.assert_called_once()
mgr._img_server_sock.close.assert_called_once()
def test_is_idempotent(self):
mgr = self._stub_with_sockets()
mgr.stop()
mgr.stop() # must not raise
def test_close_failure_is_swallowed(self):
mgr = make_manager(role=SyncRole.LEADER)
mgr._recv_sock = MagicMock()
mgr._recv_sock.close.side_effect = OSError("already closed")
mgr.stop() # must not raise
assert mgr.logger.debug.called
def test_handles_unset_sockets(self):
make_manager(role=SyncRole.STANDALONE).stop() # all sockets None
class TestFollowerAnnounceLoop:
"""The follower's outbound half of the handshake.
Covered on mock sockets so it does not depend on the network
delivering anything: the real-socket handshake below skips when the
environment drops broadcast, and that skip is only safe because a
regression in what the follower *sends* is caught here instead.
"""
def _run_once(self, monkeypatch, mgr, now=1000.0):
fake_clock(monkeypatch, time_fn=lambda: now,
sleep_fn=lambda _: setattr(mgr, "_running", False))
mgr._running = True
mgr._follower_announce_loop()
def _sent(self, mgr):
return [(json.loads(payload.decode("utf-8")), dest)
for payload, dest in
(call[0] for call in mgr._send_sock.sendto.call_args_list)]
def test_hello_carries_this_display_and_goes_to_broadcast(self, monkeypatch):
mgr = make_manager(role=SyncRole.FOLLOWER,
hw_config={"rows": 64, "cols": 128, "chain_length": 3})
mgr._send_sock = MagicMock()
self._run_once(monkeypatch, mgr)
sent = self._sent(mgr)
assert all(dest == ("<broadcast>", mgr.port) for _, dest in sent)
assert {"t": "hello", "rows": 64, "cols": 128, "chain": 3} in [m for m, _ in sent]
def test_heartbeat_is_announced_too(self, monkeypatch):
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._send_sock = MagicMock()
self._run_once(monkeypatch, mgr)
assert {"t": "hb"} in [m for m, _ in self._sent(mgr)]
def test_hello_defaults_when_hardware_config_is_empty(self, monkeypatch):
mgr = make_manager(role=SyncRole.FOLLOWER, hw_config={})
mgr._send_sock = MagicMock()
self._run_once(monkeypatch, mgr)
hello = next(m for m, _ in self._sent(mgr) if m["t"] == "hello")
assert (hello["rows"], hello["cols"], hello["chain"]) == (32, 64, 1)
def test_hello_is_not_resent_before_its_interval(self, monkeypatch):
# Heartbeat is the faster of the two, so advancing by one heartbeat
# per iteration must produce more heartbeats than hellos.
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._send_sock = MagicMock()
clock = {"now": 1000.0}
ticks = {"n": 0}
def tick(_):
ticks["n"] += 1
clock["now"] += sync_manager.HEARTBEAT_INTERVAL
if ticks["n"] >= 2:
mgr._running = False
fake_clock(monkeypatch, time_fn=lambda: clock["now"], sleep_fn=tick)
mgr._running = True
mgr._follower_announce_loop()
kinds = [m["t"] for m, _ in self._sent(mgr)]
assert kinds.count("hello") == 1
assert kinds.count("hb") == 2
def test_send_failure_is_swallowed(self, monkeypatch):
# This swallow is why a network that drops broadcast looks like
# silence rather than an error — the handshake test's skip exists
# for exactly that reason.
mgr = make_manager(role=SyncRole.FOLLOWER)
mgr._send_sock = MagicMock()
mgr._send_sock.sendto.side_effect = OSError("network unreachable")
self._run_once(monkeypatch, mgr) # must not raise
assert mgr.logger.debug.called
def _broadcast_available(port):
"""True when a UDP broadcast can be sent at all in this environment.
The handshake below depends on broadcast: the follower announces
itself to ("<broadcast>", port), and sync_manager swallows any sendto
error. Without this probe, a sandbox or CI network that refuses
broadcast would make the test wait out its whole deadline and then
fail for a reason that has nothing to do with the code.
This catches only refusal, not silent drop — confirming delivery
would mean binding INADDR_ANY to receive, a listening socket this
suite has no business opening. The drop case is handled at the
deadline instead; see the skip in the handshake test.
"""
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
try:
sock.setsockopt(socket.SOL_SOCKET, socket.SO_BROADCAST, 1)
sock.sendto(b"probe", ("<broadcast>", port))
return True
except OSError:
return False
finally:
sock.close()
class TestRealSocketHandshake:
def test_leader_and_follower_negotiate_over_real_sockets(self, monkeypatch):
# One end-to-end check that the wire format actually round-trips:
# every other test drives the loops with mocked sockets.
#
# Not loopback-only, despite the free-port probe below: the manager
# binds UDP and TCP on all interfaces and the follower announces by
# broadcast. That is the behaviour under test, so the environment
# has to support it.
monkeypatch.setattr(sync_manager, "HELLO_INTERVAL", 0.02)
monkeypatch.setattr(sync_manager, "HEARTBEAT_INTERVAL", 0.02)
hw = {"rows": 32, "cols": 64, "chain_length": 1}
leader = follower = None
# The free-port probe is inherently racy — the port can be taken
# between release and rebind — so retry rather than fail on it.
for _attempt in range(5):
# Probed on loopback: this only needs a port number, and the
# manager's own bind is what has to succeed. If the port turns
# out to be taken on another interface, the retry below covers
# it — same as for the race.
probe = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
probe.bind(("127.0.0.1", 0))
port = probe.getsockname()[1]
probe.close()
if not _broadcast_available(port):
pytest.skip("environment refuses UDP broadcast")
try:
leader = DisplaySyncManager("leader", {"port": port}, hw, MagicMock())
follower = DisplaySyncManager("follower", {"port": port}, hw, MagicMock())
break
except OSError:
# Port taken between probe and bind, or the TCP image
# server could not bind port+1. Tear down whichever end
# came up before retrying with a fresh port.
for mgr in (leader, follower):
if mgr is not None:
mgr.stop()
leader = follower = None
else:
pytest.skip("could not obtain a free port pair for the handshake")
try:
deadline = time.time() + 5.0
while time.time() < deadline:
if (leader._leader_state is LeaderState.CONNECTED
and follower._peer_compatible):
break
time.sleep(0.02)
if (leader._leader_state is LeaderState.NO_PEER
and follower._leader_ip is None):
# Not one packet crossed, in either direction. The sendto
# succeeded — _broadcast_available checked — so this is a
# network that accepts a broadcast and drops it, which no
# up-front probe can detect without binding INADDR_ANY to
# listen for its own datagram. Skip rather than report a
# protocol failure the code did not cause.
#
# This cannot hide a real regression in the announcing
# side: TestFollowerAnnounceLoop covers that on mock
# sockets, where delivery is not a variable.
pytest.skip(
"environment accepted the broadcast but did not deliver it")
assert leader._leader_state is LeaderState.CONNECTED
assert follower._peer_compatible is True
assert follower._leader_ip is not None
finally:
leader.stop()
follower.stop()