diff --git a/qtoggleserver/core/events/base.py b/qtoggleserver/core/events/base.py index 8b9d411d..d77d7097 100644 --- a/qtoggleserver/core/events/base.py +++ b/qtoggleserver/core/events/base.py @@ -14,7 +14,6 @@ class Event(metaclass=abc.ABCMeta): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_NONE TYPE = "base-event" - _UNINITIALIZED: dict = {} def __init__(self, timestamp: float | None = None) -> None: self._type: str = self.TYPE @@ -25,28 +24,29 @@ def __init__(self, timestamp: float | None = None) -> None: timestamp = 0 self._timestamp: float = timestamp - self._params: GenericJSONDict | None = self._UNINITIALIZED + self._params: GenericJSONDict | None = None def __str__(self) -> str: return f"{self._type} event" async def to_json(self) -> GenericJSONDict: - if self._params is self._UNINITIALIZED: - raise Exception("Parameters are uninitialized") - - result: GenericJSONDict = {"type": self._type} - - if self._params: - result["params"] = self._params - - return result + return { + "type": self._type, + "params": self.get_params(), + } async def init_params(self) -> None: - self._params = await self.get_params() + self._params = await self.make_params() - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return {} + def get_params(self) -> GenericJSONDict: + if self._params is None: + raise RuntimeError("Parameters are uninitialized") + + return self._params + def get_type(self) -> str: return self._type diff --git a/qtoggleserver/core/events/device.py b/qtoggleserver/core/events/device.py index a756b6e1..041e0f02 100644 --- a/qtoggleserver/core/events/device.py +++ b/qtoggleserver/core/events/device.py @@ -14,7 +14,7 @@ class DeviceUpdate(DeviceEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_ADMIN TYPE = "device-update" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return await self.get_attrs() def is_duplicate(self, event: Event) -> bool: diff --git a/qtoggleserver/core/events/port.py b/qtoggleserver/core/events/port.py index 9088ea02..79275f69 100644 --- a/qtoggleserver/core/events/port.py +++ b/qtoggleserver/core/events/port.py @@ -27,7 +27,7 @@ class PortAdd(PortEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_VIEWONLY TYPE = "port-add" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return await self.get_port().to_json() @@ -35,7 +35,7 @@ class PortRemove(PortEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_VIEWONLY TYPE = "port-remove" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return {"id": self.get_port().get_id()} @@ -43,7 +43,7 @@ class PortUpdate(PortEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_VIEWONLY TYPE = "port-update" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return await self.get_port().to_json() def is_duplicate(self, event: Event) -> bool: @@ -60,5 +60,5 @@ def __init__(self, old_value: NullablePortValue, new_value: NullablePortValue, * super().__init__(*args, **kwargs) - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return {"id": self.get_port().get_id(), "value": self.new_value, "old_value": self.old_value} diff --git a/qtoggleserver/frontend/events.py b/qtoggleserver/frontend/events.py index 7ed4032d..03181a05 100644 --- a/qtoggleserver/frontend/events.py +++ b/qtoggleserver/frontend/events.py @@ -30,5 +30,5 @@ def __init__(self, panels: GenericJSONList, **kwargs) -> None: super().__init__(**kwargs) - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return {"panels": self.panels} diff --git a/qtoggleserver/peripherals/__init__.py b/qtoggleserver/peripherals/__init__.py index 8ba3cb31..89c8aa4b 100644 --- a/qtoggleserver/peripherals/__init__.py +++ b/qtoggleserver/peripherals/__init__.py @@ -1,4 +1,5 @@ import asyncio +import hashlib import logging from collections.abc import ValuesView @@ -6,6 +7,7 @@ from qtoggleserver import persist from qtoggleserver.conf import settings +from qtoggleserver.core.ports import BasePort from qtoggleserver.utils import dynload as dynload_utils from .exceptions import DuplicatePeripheral, NoSuchDriver @@ -91,9 +93,6 @@ async def prepare_migration( - Phase 1 (prepare): Copy data to new IDs, keep old data intact - Phase 2 (cleanup): Delete old data only after successful update """ - import hashlib - - from qtoggleserver.core.ports import BasePort # Compute what the new peripheral ID will be using the same logic as Peripheral.__init__ new_id: str = new_name or "" @@ -141,7 +140,6 @@ async def cleanup_migration(p: Peripheral, new_name: str | None) -> None: Deletes old port persist data and old peripheral persist entry. Only call this after the new peripheral has been successfully created and initialized. """ - from qtoggleserver.core.ports import BasePort old_id = p.get_id() diff --git a/qtoggleserver/peripherals/events.py b/qtoggleserver/peripherals/events.py index 8acc079c..65f2e57d 100644 --- a/qtoggleserver/peripherals/events.py +++ b/qtoggleserver/peripherals/events.py @@ -22,7 +22,7 @@ class PeripheralAdd(PeripheralEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_ADMIN TYPE = "peripheral-add" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return self.get_peripheral().to_json() @@ -30,7 +30,7 @@ class PeripheralRemove(PeripheralEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_ADMIN TYPE = "peripheral-remove" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return {"id": self.get_peripheral().get_id()} @@ -38,7 +38,7 @@ class PeripheralUpdate(PeripheralEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_ADMIN TYPE = "peripheral-update" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return self.get_peripheral().to_json() def is_duplicate(self, event: core_events.Event) -> bool: diff --git a/qtoggleserver/slaves/devices.py b/qtoggleserver/slaves/devices.py index 15f50190..227d314a 100644 --- a/qtoggleserver/slaves/devices.py +++ b/qtoggleserver/slaves/devices.py @@ -1199,6 +1199,7 @@ async def _handle_offline(self) -> None: # Trigger a port-update so that online attribute is pushed to consumers for port in self._get_local_ports(): if port.is_enabled(): + port.invalidate_attrs() await port.trigger_update() async def _handle_online(self) -> None: @@ -1232,6 +1233,7 @@ async def _handle_online(self) -> None: # Trigger a port-update so that online attribute is pushed to consumers for port in self._get_local_ports(): if port.is_enabled(): + port.invalidate_attrs() await port.trigger_update() if not self._ready: diff --git a/qtoggleserver/slaves/events.py b/qtoggleserver/slaves/events.py index 80ea4299..35efb63e 100644 --- a/qtoggleserver/slaves/events.py +++ b/qtoggleserver/slaves/events.py @@ -26,7 +26,7 @@ class SlaveDeviceAdd(SlaveDeviceEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_ADMIN TYPE = "slave-device-add" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return self.get_slave().to_json() @@ -34,7 +34,7 @@ class SlaveDeviceRemove(SlaveDeviceEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_ADMIN TYPE = "slave-device-remove" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return {"name": self.get_slave().get_name()} @@ -42,7 +42,7 @@ class SlaveDeviceUpdate(SlaveDeviceEvent): REQUIRED_ACCESS = core_api.ACCESS_LEVEL_ADMIN TYPE = "slave-device-update" - async def get_params(self) -> GenericJSONDict: + async def make_params(self) -> GenericJSONDict: return self.get_slave().to_json() def is_duplicate(self, event: core_events.Event) -> bool: diff --git a/qtoggleserver/slaves/ports.py b/qtoggleserver/slaves/ports.py index 7df0d1e9..9d7d3922 100644 --- a/qtoggleserver/slaves/ports.py +++ b/qtoggleserver/slaves/ports.py @@ -355,6 +355,7 @@ def heart_beat_second(self) -> None: self.debug("value expired") if not self._trigger_update_task: + self.invalidate_attrs() self._trigger_update_task = asyncio.create_task(self.trigger_update()) async def from_persisted(self, data: GenericJSONDict) -> None: diff --git a/qtoggleserver/startup.py b/qtoggleserver/startup.py index 519b8ccd..98c1b2a9 100644 --- a/qtoggleserver/startup.py +++ b/qtoggleserver/startup.py @@ -323,19 +323,6 @@ async def init_main() -> None: logger.info("initializing main") await main.init() - # Wait until slaves are also ready before actually considering main loop ready - if settings.slaves.enabled: - logger.debug("waiting for slaves to become ready") - while not slaves_devices.ready(): - await asyncio.sleep(1) - - logger.debug("slaves are ready") - - # Mark main as ready after all slaves with their ports have been initialized and hopefully brought online. Allow an - # extra second for pending loop tasks. - await asyncio.sleep(1) - main.set_ready() - async def cleanup_main() -> None: logger.info("cleaning up main") @@ -370,11 +357,24 @@ async def init() -> None: await init_device() await init_webhooks() await init_reverse() + await init_main() await init_ports() await init_slaves() - await init_main() await init_web() + # Wait until slaves are also ready before actually considering main loop ready + if settings.slaves.enabled: + logger.debug("waiting for slaves to become ready") + while not slaves_devices.ready(): + await asyncio.sleep(1) + + logger.debug("slaves are ready") + + # Mark main as ready after all slaves with their ports have been initialized and hopefully brought online. Allow an + # extra second for pending loop tasks. + await asyncio.sleep(1) + main.set_ready() + async def cleanup() -> None: await cleanup_web() diff --git a/qtoggleserver/utils/expressions.py b/qtoggleserver/utils/expressions.py index 4f50a6c1..b87b2013 100644 --- a/qtoggleserver/utils/expressions.py +++ b/qtoggleserver/utils/expressions.py @@ -38,9 +38,10 @@ def invalidate_deps_map() -> None: async def build_context(now_ms: int) -> EvalContext: """Build an expression evaluation context for the current system state. - Gathers port values and attributes for all enabled ports, collects device-level attributes, - and includes slave device attributes if slaves are enabled. Returns a complete EvalContext - ready for expression evaluation. + Gathers port values and attributes for all ports (including disabled ones; disabled ports are + only special-cased for port *value* expressions, not attribute expressions), collects + device-level attributes, and includes slave device attributes if slaves are enabled. Returns a + complete EvalContext ready for expression evaluation. Args: now_ms: Current time in milliseconds since epoch. @@ -51,9 +52,6 @@ async def build_context(now_ms: int) -> EvalContext: port_values = {} port_attrs = {} for port in core_ports.get_all(): - if not port.is_enabled(): - continue - port_id = port.get_id() port_values[port_id] = port.get_last_value() port_attrs[port_id] = await port.get_attrs() diff --git a/qtoggleserver/utils/main.py b/qtoggleserver/utils/main.py index a61609af..412db9a5 100644 --- a/qtoggleserver/utils/main.py +++ b/qtoggleserver/utils/main.py @@ -11,15 +11,24 @@ class AttrChangeHandler(core_events.Handler): Dep strings produced: - ``$port_id:`` — for "port-add", "port-remove", "port-update" + - ``$port_id`` — additionally, for "port-add"/"port-remove"/"port-update", whenever a port's availability + (`enabled`/`online`) transitions, in either direction, relative to its last known state — so that value + expressions and functions like `AVAILABLE()`/`DEFAULT()` pick up the new state - ``#:`` — for "device-update" - ``#name:`` — for "slave-device-add", "slave-device-remove", "slave-device-update" """ FIRE_AND_FORGET = False + # Attribute names whose transition (in either direction) changes whether a port's *value* is available, and + # therefore must also be treated as a change of the port's value dependency (`$port_id`), not just its attribute + # dependency (`$port_id:`). + AVAILABILITY_ATTRS = ("enabled", "online") + def __init__(self) -> None: super().__init__(name="attribute-changes") self._pending: set[str] = set() + self._last_availability: dict[tuple[str, str], bool] = {} def pop_pending(self) -> set[str]: """Return the pending changes and clear the internal set.""" @@ -29,7 +38,26 @@ def pop_pending(self) -> set[str]: async def handle_event(self, event: core_events.Event) -> None: if isinstance(event, (core_events.PortAdd, core_events.PortRemove, core_events.PortUpdate)): - self._pending.add(f"${event.get_port().get_id()}:") + port_id = event.get_port().get_id() + self._pending.add(f"${port_id}:") + + if isinstance(event, core_events.PortRemove): + for attr_name in self.AVAILABILITY_ATTRS: + # Whenever one of the availability attributes changes, induce a `value-change`-like event so that + # expressions depending on this port's value are re-evaluated. + was_available = self._last_availability.pop((port_id, attr_name), False) + if was_available: + self._pending.add(f"${port_id}") + else: + # Also covers "port-add" — see class docstring. + params = event.get_params() + for attr_name in self.AVAILABILITY_ATTRS: + # Whenever one of the availability attributes changes, induce a `value-change`-like event so that + # expressions depending on this port's value are re-evaluated. + available = bool(params.get(attr_name)) + if available != self._last_availability.get((port_id, attr_name), False): + self._pending.add(f"${port_id}") + self._last_availability[(port_id, attr_name)] = available elif isinstance(event, core_events.DeviceUpdate): self._pending.add("#:") elif isinstance( diff --git a/tests/unit/qtoggleserver/core/test_main.py b/tests/unit/qtoggleserver/core/test_main.py index 8466181a..1d001918 100644 --- a/tests/unit/qtoggleserver/core/test_main.py +++ b/tests/unit/qtoggleserver/core/test_main.py @@ -512,54 +512,268 @@ async def test_pause_resume_cycle(self, freezer, mocker, mock_num_port1, dummy_u class TestAttrChangeHandler: @pytest.fixture(autouse=True) def reset_pending(self): - """Clear pending attr changes before and after each test.""" + """Clear pending attr changes and availability-tracking state before and after each test.""" core_main._attr_change_handler._pending.clear() + core_main._attr_change_handler._last_availability.clear() yield core_main._attr_change_handler._pending.clear() + core_main._attr_change_handler._last_availability.clear() async def test_port_update_adds_attr_dep(self, mock_num_port1): - """PortUpdate event should add `$port_id:` to _pending_attr_changes.""" + """PortUpdate event should add `$port_id:` (and `$port_id`, since the mock port is already enabled) to + _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(core_events.PortUpdate(mock_num_port1)) - assert core_main._attr_change_handler._pending == {"$nid1:"} + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} async def test_port_add_adds_attr_dep(self, mock_num_port1): - """PortAdd event should add `$port_id:` to _pending_attr_changes.""" + """PortAdd event should add `$port_id:` (and `$port_id`, since the mock port is already enabled) to + _pending_attr_changes.""" + + event = core_events.PortAdd(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} + + async def test_port_add_while_disabled_does_not_add_value_dep(self, mocker, mock_num_port1): + """PortAdd should not add `$port_id` when the port is added while already disabled — a newly added disabled + port is not a real availability transition, since a nonexistent port was already unavailable.""" + + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": False}) + event = core_events.PortAdd(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:"} + + async def test_port_add_seeds_availability_so_later_update_does_not_repeat(self, mock_num_port1): + """After a PortAdd for an already-enabled port seeds `_last_availability`, a later unrelated PortUpdate (with + `enabled` still `true`) should not add `$port_id` again — this is what makes registering the handler before + ports are loaded at startup (see startup.py) actually pay off, instead of every port's first post-startup + event looking like a spurious availability transition.""" - await core_main._attr_change_handler.handle_event(core_events.PortAdd(mock_num_port1)) + add_event = core_events.PortAdd(mock_num_port1) + await add_event.init_params() + await core_main._attr_change_handler.handle_event(add_event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} + core_main._attr_change_handler._pending.clear() + + update_event = core_events.PortUpdate(mock_num_port1) + await update_event.init_params() + await core_main._attr_change_handler.handle_event(update_event) assert core_main._attr_change_handler._pending == {"$nid1:"} async def test_port_remove_adds_attr_dep(self, mock_num_port1): """PortRemove event should add `$port_id:` to _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(core_events.PortRemove(mock_num_port1)) + event = core_events.PortRemove(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:"} + + async def test_enabled_transition_adds_value_dep(self, mocker, mock_num_port1): + """`$port_id` should be added whenever `enabled` is observed transitioning, in either direction.""" + + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": False}) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:"} + core_main._attr_change_handler._pending.clear() + + mock_num_port1.to_json.return_value = {"enabled": True} + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} # rising edge + core_main._attr_change_handler._pending.clear() + + mock_num_port1.to_json.return_value = {"enabled": False} + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} # falling edge + + async def test_enabled_steady_state_does_not_repeat(self, mocker, mock_num_port1): + """`$port_id` should only be added once while `enabled` stays `true` across events.""" + + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": True}) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler._pending.clear() + + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:"} + + async def test_disabled_steady_state_does_not_repeat(self, mocker, mock_num_port1): + """`$port_id` should not be added again while `enabled` stays `false` across events.""" + + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": False}) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler._pending.clear() + + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) assert core_main._attr_change_handler._pending == {"$nid1:"} + async def test_online_transition_adds_value_dep(self, mocker, mock_num_port1): + """`$port_id` should be added when `online` (independently of `enabled`) transitions to `true`.""" + + mocker.patch.object( + mock_num_port1, + "to_json", + new_callable=mocker.AsyncMock, + return_value={"enabled": True, "online": False}, + ) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} # from the `enabled` transition + core_main._attr_change_handler._pending.clear() + + mock_num_port1.to_json.return_value = {"enabled": True, "online": True} + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} # from the `online` transition + core_main._attr_change_handler._pending.clear() + + mock_num_port1.to_json.return_value = {"enabled": True, "online": False} + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} # from the `online` falling edge + + async def test_enabled_and_online_same_event_adds_dep_once(self, mocker, mock_num_port1): + """`$port_id` should be added only once even if both `enabled` and `online` newly become `true`.""" + + mocker.patch.object( + mock_num_port1, + "to_json", + new_callable=mocker.AsyncMock, + return_value={"enabled": True, "online": True}, + ) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} + + async def test_enabled_and_online_same_event_removes_dep_once(self, mocker, mock_num_port1): + """`$port_id` should be added only once even if both `enabled` and `online` newly become `false`.""" + + mocker.patch.object( + mock_num_port1, + "to_json", + new_callable=mocker.AsyncMock, + return_value={"enabled": True, "online": True}, + ) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler._pending.clear() + + mock_num_port1.to_json.return_value = {"enabled": False, "online": False} + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} + + async def test_port_remove_while_available_adds_value_dep(self, mocker, mock_num_port1): + """PortRemove should add `$port_id` too when the port was available right before being removed, since + removal makes its value unavailable — this is what makes `AVAILABLE()`/`DEFAULT()` pick up the removal.""" + + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": True}) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler._pending.clear() + + event = core_events.PortRemove(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} + + async def test_port_remove_while_unavailable_does_not_add_value_dep(self, mocker, mock_num_port1): + """PortRemove should not add `$port_id` when the port was already unavailable before being removed.""" + + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": False}) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler._pending.clear() + + event = core_events.PortRemove(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:"} + + async def test_port_remove_clears_availability_tracking(self, mocker, mock_num_port1): + """PortRemove should clear tracked availability so a later re-add is treated as a fresh transition.""" + + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": True}) + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler._pending.clear() + + event = core_events.PortRemove(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler._pending.clear() + + event = core_events.PortAdd(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1"} + async def test_device_update_adds_device_dep(self): """DeviceUpdate event should add `#:` to _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(core_events.DeviceUpdate()) + event = core_events.DeviceUpdate() + await event.init_params() + await core_main._attr_change_handler.handle_event(event) assert core_main._attr_change_handler._pending == {"#:"} async def test_value_change_ignored(self, mock_num_port1): """ValueChange event should not modify _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(core_events.ValueChange(None, 42, mock_num_port1)) + event = core_events.ValueChange(None, 42, mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) assert core_main._attr_change_handler._pending == set() async def test_multiple_events_accumulate(self, mock_num_port1, mock_num_port2): """Multiple port events should accumulate their dep strings.""" - await core_main._attr_change_handler.handle_event(core_events.PortUpdate(mock_num_port1)) - await core_main._attr_change_handler.handle_event(core_events.PortUpdate(mock_num_port2)) - assert core_main._attr_change_handler._pending == {"$nid1:", "$nid2:"} + event1 = core_events.PortUpdate(mock_num_port1) + await event1.init_params() + await core_main._attr_change_handler.handle_event(event1) + + event2 = core_events.PortUpdate(mock_num_port2) + await event2.init_params() + await core_main._attr_change_handler.handle_event(event2) + + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1", "$nid2:", "$nid2"} async def test_port_and_device_events_accumulate(self, mock_num_port1): """Port and device events should both accumulate into _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(core_events.PortUpdate(mock_num_port1)) - await core_main._attr_change_handler.handle_event(core_events.DeviceUpdate()) - assert core_main._attr_change_handler._pending == {"$nid1:", "#:"} + port_event = core_events.PortUpdate(mock_num_port1) + await port_event.init_params() + await core_main._attr_change_handler.handle_event(port_event) + + device_event = core_events.DeviceUpdate() + await device_event.init_params() + await core_main._attr_change_handler.handle_event(device_event) + + assert core_main._attr_change_handler._pending == {"$nid1:", "$nid1", "#:"} async def test_read_ports_drains_pending_attr_changes(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): """read_ports() should include pending port-attr deps in the changes set passed to _eval_changed_expressions.""" @@ -610,6 +824,76 @@ async def test_port_attr_dep_triggers_expression_eval(self, mocker, mock_num_por mock_num_port2.eval_and_push_write.assert_called_once() + async def test_port_enabled_transition_triggers_dependent_value_eval(self, mocker, mock_num_port1, mock_num_port2): + """A port whose expression references $nid1 (value dep) should be evaluated once nid1's `enabled` transition + is observed via a PortUpdate event and drained into the deps map lookup, end to end.""" + + mock_num_port2.set_writable(True) + mock_num_port2.set_expression("$nid1") + mocker.patch.object(mock_num_port2, "eval_and_push_write") + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": True}) + + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + changes = core_main._attr_change_handler.pop_pending() + assert "$nid1" in changes + + await _eval_changed_expressions(changes=changes, now_ms=0) + + mock_num_port2.eval_and_push_write.assert_called_once() + + async def test_port_disabled_transition_triggers_dependent_availability_eval( + self, mocker, mock_num_port1, mock_num_port2 + ): + """A port whose expression references AVAILABLE($nid1) should be evaluated once nid1's `enabled` transitions + to `false`, so that it picks up the port becoming unavailable, end to end.""" + + mock_num_port2.set_writable(True) + mock_num_port2.set_expression("AVAILABLE($nid1)") + mocker.patch.object(mock_num_port2, "eval_and_push_write") + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": True}) + + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler.pop_pending() # discard the initial rising-edge dep + + mock_num_port1.to_json.return_value = {"enabled": False} + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + changes = core_main._attr_change_handler.pop_pending() + assert "$nid1" in changes + + await _eval_changed_expressions(changes=changes, now_ms=0) + + mock_num_port2.eval_and_push_write.assert_called_once() + + async def test_port_remove_triggers_dependent_default_eval(self, mocker, mock_num_port1, mock_num_port2): + """A port whose expression references DEFAULT($nid1, ...) should be evaluated once nid1 is removed while + available, so that it picks up the port becoming unavailable, end to end.""" + + mock_num_port2.set_writable(True) + mock_num_port2.set_expression("DEFAULT($nid1, 0)") + mocker.patch.object(mock_num_port2, "eval_and_push_write") + mocker.patch.object(mock_num_port1, "to_json", new_callable=mocker.AsyncMock, return_value={"enabled": True}) + + event = core_events.PortUpdate(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + core_main._attr_change_handler.pop_pending() # discard the initial rising-edge dep + + event = core_events.PortRemove(mock_num_port1) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) + changes = core_main._attr_change_handler.pop_pending() + assert "$nid1" in changes + + await _eval_changed_expressions(changes=changes, now_ms=0) + + mock_num_port2.eval_and_push_write.assert_called_once() + async def test_device_dep_triggers_expression_eval(self, mocker, mock_num_port1): """A port whose expression references #:attr should be evaluated when `#:` is in changes.""" @@ -624,26 +908,38 @@ async def test_device_dep_triggers_expression_eval(self, mocker, mock_num_port1) async def test_slave_device_update_adds_slave_dep(self): """SlaveDeviceUpdate event should add `#slave_name:` to _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1"))) + event = slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1")) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) assert core_main._attr_change_handler._pending == {"#slave1:"} async def test_slave_device_add_adds_slave_dep(self): """SlaveDeviceAdd event should add `#slave_name:` to _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceAdd(_make_mock_slave("slave1"))) + event = slaves_events.SlaveDeviceAdd(_make_mock_slave("slave1")) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) assert core_main._attr_change_handler._pending == {"#slave1:"} async def test_slave_device_remove_adds_slave_dep(self): """SlaveDeviceRemove event should add `#slave_name:` to _pending_attr_changes.""" - await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceRemove(_make_mock_slave("slave1"))) + event = slaves_events.SlaveDeviceRemove(_make_mock_slave("slave1")) + await event.init_params() + await core_main._attr_change_handler.handle_event(event) assert core_main._attr_change_handler._pending == {"#slave1:"} async def test_multiple_slave_events_accumulate(self): """Multiple slave device events should accumulate distinct dep strings.""" - await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1"))) - await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave2"))) + event1 = slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1")) + await event1.init_params() + await core_main._attr_change_handler.handle_event(event1) + + event2 = slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave2")) + await event2.init_params() + await core_main._attr_change_handler.handle_event(event2) + assert core_main._attr_change_handler._pending == {"#slave1:", "#slave2:"} async def test_slave_dep_triggers_expression_eval(self, mocker, mock_num_port1): diff --git a/tests/unit/qtoggleserver/peripherals/test_events.py b/tests/unit/qtoggleserver/peripherals/test_events.py index d4566373..0b2d8776 100644 --- a/tests/unit/qtoggleserver/peripherals/test_events.py +++ b/tests/unit/qtoggleserver/peripherals/test_events.py @@ -4,31 +4,31 @@ class TestPeripheralAdd: - async def test_get_params(self): + async def test_make_params(self): peripheral = MockPeripheral(name="test", dummy_param="v") event = peripherals_events.PeripheralAdd(peripheral) - assert await event.get_params() == peripheral.to_json() + assert await event.make_params() == peripheral.to_json() assert event.TYPE == "peripheral-add" assert event.REQUIRED_ACCESS == core_api.ACCESS_LEVEL_ADMIN class TestPeripheralRemove: - async def test_get_params(self): + async def test_make_params(self): peripheral = MockPeripheral(name="test", dummy_param="v") event = peripherals_events.PeripheralRemove(peripheral) - assert await event.get_params() == {"id": peripheral.get_id()} + assert await event.make_params() == {"id": peripheral.get_id()} assert event.TYPE == "peripheral-remove" assert event.REQUIRED_ACCESS == core_api.ACCESS_LEVEL_ADMIN class TestPeripheralUpdate: - async def test_get_params(self): + async def test_make_params(self): peripheral = MockPeripheral(name="test", dummy_param="v") event = peripherals_events.PeripheralUpdate(peripheral) - assert await event.get_params() == peripheral.to_json() + assert await event.make_params() == peripheral.to_json() assert event.TYPE == "peripheral-update" assert event.REQUIRED_ACCESS == core_api.ACCESS_LEVEL_ADMIN diff --git a/tests/unit/qtoggleserver/test_startup.py b/tests/unit/qtoggleserver/test_startup.py index e3ffd1a1..86d19c50 100644 --- a/tests/unit/qtoggleserver/test_startup.py +++ b/tests/unit/qtoggleserver/test_startup.py @@ -80,6 +80,77 @@ async def test_continues_when_static_port_load_has_partial_errors(self, mocker): assert spy_logger.error.call_count == 2 +class TestInit: + _INIT_STEPS = [ + "init_metadata", + "init_system", + "init_persist", + "init_peripherals", + "init_events", + "init_sessions", + "init_history", + "init_device", + "init_webhooks", + "init_reverse", + "init_main", + "init_ports", + "init_slaves", + "init_web", + ] + + def _patch_steps(self, mocker, call_order=None): + for name in self._INIT_STEPS: + side_effect = (lambda n=name: call_order.append(n)) if call_order is not None else None + mocker.patch(f"qtoggleserver.startup.{name}", side_effect=side_effect) + + mocker.patch("qtoggleserver.startup.parse_args") + mocker.patch("qtoggleserver.startup.init_settings") + mocker.patch("qtoggleserver.startup.init_logging") + mocker.patch("qtoggleserver.startup.init_signals") + mocker.patch("qtoggleserver.startup.init_tornado") + mocker.patch("qtoggleserver.startup.logger", mocker.MagicMock()) + mocker.patch.object(startup.settings.slaves, "enabled", False) + mocker.patch("asyncio.sleep", new_callable=mocker.AsyncMock) + + async def test_init_main_runs_before_ports_and_slaves(self, mocker): + """init_main() (which registers core.main's attr-change event handler) must run before init_ports() and + init_slaves(), so that PortAdd events fired while loading pre-existing ports are actually observed by the + handler instead of being fired into an empty handler list.""" + + call_order = [] + self._patch_steps(mocker, call_order) + mocker.patch("qtoggleserver.startup.main.set_ready") + + await startup.init() + + assert call_order.index("init_main") < call_order.index("init_ports") + assert call_order.index("init_main") < call_order.index("init_slaves") + + async def test_marks_main_ready_after_full_init(self, mocker): + """init() should call main.set_ready() once initialization (including ports/slaves) has completed.""" + + self._patch_steps(mocker) + spy_set_ready = mocker.patch("qtoggleserver.startup.main.set_ready") + + await startup.init() + + spy_set_ready.assert_called_once_with() + + async def test_waits_for_slaves_ready_before_marking_main_ready(self, mocker): + """init() should poll slaves_devices.ready() until true, when slaves are enabled, before calling + main.set_ready().""" + + self._patch_steps(mocker) + mocker.patch.object(startup.settings.slaves, "enabled", True) + spy_set_ready = mocker.patch("qtoggleserver.startup.main.set_ready") + spy_ready = mocker.patch("qtoggleserver.startup.slaves_devices.ready", side_effect=[False, True]) + + await startup.init() + + assert spy_ready.call_count == 2 + spy_set_ready.assert_called_once_with() + + class TestInitPeripherals: async def test_triggers_peripheral_add_events(self, mocker): peripheral1 = mocker.MagicMock() diff --git a/tests/unit/qtoggleserver/utils/test_expressions.py b/tests/unit/qtoggleserver/utils/test_expressions.py index 1875f7e7..0c5bcc5f 100644 --- a/tests/unit/qtoggleserver/utils/test_expressions.py +++ b/tests/unit/qtoggleserver/utils/test_expressions.py @@ -86,14 +86,12 @@ async def test_expression_cleared_excluded_from_rebuilt_map(self, mock_num_port1 class TestBuildContext: async def test_basic_context(self, mock_num_port1, mock_num_port2, mocker): - """Should gather port values and attributes from all enabled ports and create EvalContext.""" + """Should gather port values and attributes from all ports and create EvalContext.""" mocker.patch( "qtoggleserver.utils.expressions.core_ports.get_all", return_value=[mock_num_port1, mock_num_port2], ) - mocker.patch.object(mock_num_port1, "is_enabled", return_value=True) - mocker.patch.object(mock_num_port2, "is_enabled", return_value=True) mocker.patch.object(mock_num_port1, "get_id", return_value="nid1") mocker.patch.object(mock_num_port2, "get_id", return_value="nid2") mocker.patch.object(mock_num_port1, "get_last_value", return_value=42) @@ -116,8 +114,14 @@ async def test_basic_context(self, mock_num_port1, mock_num_port2, mocker): assert context.now_ms == 1000 assert context.timestamp == 1 - async def test_disabled_port_excluded(self, mock_num_port1, mock_num_port2, mocker): - """Should exclude disabled ports from the context.""" + async def test_disabled_port_included(self, mock_num_port1, mock_num_port2, mocker): + """Should still include a disabled port's value and attrs in the context. + + Per spec, only *value* expressions special-case disabled ports (evaluating to + unavailable); attribute expressions on a disabled port must still resolve normally. That + distinction is enforced by `PortValue`/`PortAttr` eval logic, not by `build_context`, so + the context must contain data for disabled ports too. + """ mocker.patch( "qtoggleserver.utils.expressions.core_ports.get_all", @@ -126,8 +130,11 @@ async def test_disabled_port_excluded(self, mock_num_port1, mock_num_port2, mock mocker.patch.object(mock_num_port1, "is_enabled", return_value=True) mocker.patch.object(mock_num_port2, "is_enabled", return_value=False) mocker.patch.object(mock_num_port1, "get_id", return_value="nid1") + mocker.patch.object(mock_num_port2, "get_id", return_value="nid2") mocker.patch.object(mock_num_port1, "get_last_value", return_value=42) + mocker.patch.object(mock_num_port2, "get_last_value", return_value=84) mocker.patch.object(mock_num_port1, "get_attrs", new_callable=mocker.AsyncMock, return_value={}) + mocker.patch.object(mock_num_port2, "get_attrs", new_callable=mocker.AsyncMock, return_value={"enabled": False}) mocker.patch( "qtoggleserver.utils.expressions.core_device_attrs.get_attrs", new_callable=mocker.AsyncMock, @@ -137,8 +144,8 @@ async def test_disabled_port_excluded(self, mock_num_port1, mock_num_port2, mock context = await expressions.build_context(2000) - assert context.port_values == {"nid1": 42} - assert context.port_attrs == {"nid1": {}} + assert context.port_values == {"nid1": 42, "nid2": 84} + assert context.port_attrs == {"nid1": {}, "nid2": {"enabled": False}} assert context.now_ms == 2000 async def test_no_ports(self, mocker): @@ -170,7 +177,6 @@ async def test_with_slave_attrs(self, mock_num_port1, mocker): "qtoggleserver.utils.expressions.core_ports.get_all", return_value=[mock_num_port1], ) - mocker.patch.object(mock_num_port1, "is_enabled", return_value=True) mocker.patch.object(mock_num_port1, "get_id", return_value="nid1") mocker.patch.object(mock_num_port1, "get_last_value", return_value=42) mocker.patch.object(mock_num_port1, "get_attrs", new_callable=mocker.AsyncMock, return_value={}) @@ -205,7 +211,6 @@ async def test_with_multiple_slaves(self, mock_num_port1, mocker): "qtoggleserver.utils.expressions.core_ports.get_all", return_value=[mock_num_port1], ) - mocker.patch.object(mock_num_port1, "is_enabled", return_value=True) mocker.patch.object(mock_num_port1, "get_id", return_value="nid1") mocker.patch.object(mock_num_port1, "get_last_value", return_value=42) mocker.patch.object(mock_num_port1, "get_attrs", new_callable=mocker.AsyncMock, return_value={}) @@ -240,7 +245,6 @@ async def test_timestamp_calculation(self, mock_num_port1, mocker): "qtoggleserver.utils.expressions.core_ports.get_all", return_value=[mock_num_port1], ) - mocker.patch.object(mock_num_port1, "is_enabled", return_value=True) mocker.patch.object(mock_num_port1, "get_id", return_value="nid1") mocker.patch.object(mock_num_port1, "get_last_value", return_value=0) mocker.patch.object(mock_num_port1, "get_attrs", new_callable=mocker.AsyncMock, return_value={})