From 62115ee0ef7dd6cb2689a996a53aa17d8d1fe30d Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 4 May 2026 11:33:26 +0300 Subject: [PATCH 01/14] core/expressions: Add attributes support --- qtoggleserver/core/expressions/__init__.py | 9 ++- qtoggleserver/core/expressions/base.py | 12 ++++ qtoggleserver/core/expressions/devices.py | 68 ++++++++++++++++++ qtoggleserver/core/expressions/exceptions.py | 30 ++++++++ qtoggleserver/core/expressions/ports.py | 72 ++++++++++++++++---- 5 files changed, 175 insertions(+), 16 deletions(-) create mode 100644 qtoggleserver/core/expressions/devices.py diff --git a/qtoggleserver/core/expressions/__init__.py b/qtoggleserver/core/expressions/__init__.py index 55b00791..c5986799 100644 --- a/qtoggleserver/core/expressions/__init__.py +++ b/qtoggleserver/core/expressions/__init__.py @@ -11,6 +11,7 @@ Expression, Role, ) +from .exceptions import EmptyExpression from .literalvalues import LiteralValue @@ -52,14 +53,20 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int = 1) pos += len(sexpression) - len(stripped) sexpression = stripped.rstrip() - if sexpression and sexpression[0] in ("$", "@"): + if not sexpression: + raise EmptyExpression() + + if sexpression[0] in ("$", "@"): return PortExpression.parse(self_port_id, sexpression, role, pos) + elif sexpression[0] == "#": + return DeviceExpression.parse(self_port_id, sexpression, role, pos) elif "(" in sexpression or ")" in sexpression: return Function.parse(self_port_id, sexpression, role, pos) else: return LiteralValue.parse(self_port_id, sexpression, role, pos) +from .devices import DeviceExpression # noqa: E402 from .functions import ( # noqa: E402 Function, # noqa: E402 aggregation, diff --git a/qtoggleserver/core/expressions/base.py b/qtoggleserver/core/expressions/base.py index a5d4770f..27f86333 100644 --- a/qtoggleserver/core/expressions/base.py +++ b/qtoggleserver/core/expressions/base.py @@ -53,6 +53,18 @@ async def _eval(self, context: EvalContext) -> EvalResult: raise NotImplementedError() def get_deps(self) -> set[str]: + """ + Return a set with all dependencies of this expression. + Each dependency is a string and, depending on its format denotes a different type of dependency: + - If string starts with `$`, the dep is a port value or port attribute dependency; going further, we have: + - `$port_id` - dependency on the corresponding port's value + - `$port_id:` - dependency on the corresponding port's attributes + - If string starts with `#`, the dep is a device attribute dependency; going further, we have: + - `#:` - dependency on (main) device attributes + - `#slave_name:` - dependency on the corresponding slave's attributes + - One of the `DEP_*` constants (e.g. `DEP_YEAR` or `"year"`) indicates a time dependency + """ + if self._cached_deps is None: self._cached_deps = self._get_deps() return self._cached_deps diff --git a/qtoggleserver/core/expressions/devices.py b/qtoggleserver/core/expressions/devices.py new file mode 100644 index 00000000..9956482e --- /dev/null +++ b/qtoggleserver/core/expressions/devices.py @@ -0,0 +1,68 @@ +from __future__ import annotations + +import abc +import re + +from .base import EvalContext, EvalResult, Expression, Role +from .exceptions import DeviceAttrUnavailable, MissingAttrPrefix, UnexpectedCharacter + + +class DeviceExpression(Expression, metaclass=abc.ABCMeta): + def __init__(self, device_name: str | None, prefix: str, role: Role, attr_name: str | None = None) -> None: + super().__init__(role) + + self.device_name: str | None = device_name + self.attr_name: str | None = attr_name + self.prefix: str = prefix + + @staticmethod + def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> Expression: + stripped = sexpression.lstrip() + pos += len(sexpression) - len(stripped) + sexpression = stripped.rstrip() + + prefix = sexpression[0] + sub_sexpression = sexpression[1:] + + if prefix == "#": + parts = sub_sexpression.split(":", 1) + if len(parts) != 2: + raise MissingAttrPrefix(pos + len(sub_sexpression)) # TODO unit test + + device_name, attr_name = parts + if device_name: + m = re.search(r"[^a-zA-Z0-9_-]", device_name) + if m: + p = m.start() + raise UnexpectedCharacter(device_name[p], p + pos + 3) # TODO: is this reported position correct? + + return SlaveDeviceAttr(device_name, prefix, role, attr_name) + else: + return MainDeviceAttr(prefix, role, attr_name) + else: + raise UnexpectedCharacter(prefix, pos) + + +class DeviceAttr(DeviceExpression): + def __str__(self) -> str: + return f"{self.prefix}{self.device_name or ''}:{self.attr_name}" + + def _get_deps(self) -> set[str]: + return {f"#{self.device_name or ''}:"} + + async def _eval(self, context: EvalContext) -> EvalResult: + key = f"{self.device_name}.{self.attr_name or ''}" if self.device_name else self.attr_name + value = context.device_attrs.get(key) + if value is None: + raise DeviceAttrUnavailable(self.device_name or "", self.attr_name or "") + + return value + + +class SlaveDeviceAttr(DeviceAttr): + pass + + +class MainDeviceAttr(DeviceAttr): + def __init__(self, prefix: str, role: Role, attr_name: str | None = None) -> None: + super().__init__(None, prefix, role, attr_name) diff --git a/qtoggleserver/core/expressions/exceptions.py b/qtoggleserver/core/expressions/exceptions.py index 54470bd0..b006ee74 100644 --- a/qtoggleserver/core/expressions/exceptions.py +++ b/qtoggleserver/core/expressions/exceptions.py @@ -93,6 +93,16 @@ def to_json(self) -> GenericJSONDict: return {"reason": "empty"} +class MissingAttrPrefix(ExpressionParseError): + def __init__(self, pos: int) -> None: + self.pos: int = pos + + super().__init__("Missing attribute prefix") + + def to_json(self) -> GenericJSONDict: + return {"reason": "missing-attr-prefix", "pos": self.pos} + + class ExpressionEvalException(ExpressionException): pass @@ -127,6 +137,26 @@ class DisabledPort(PortValueUnavailable): MSG = 'Port "%s" is disabled' +class PortAttrUnavailable(ValueUnavailable): + MSG = 'Port attribute "%s:%s" is unavailable' + + def __init__(self, port_id: str, attr_name: str) -> None: + self.port_id = port_id + self.attr_name = attr_name + + super().__init__(self.MSG % (port_id, attr_name)) + + +class DeviceAttrUnavailable(ValueUnavailable): + MSG = 'Device attribute "%s:%s" is unavailable' + + def __init__(self, device_name: str, attr_name: str) -> None: + self.device_name = device_name + self.attr_name = attr_name + + super().__init__(self.MSG % (device_name, attr_name)) + + class ExpressionArithmeticError(ExpressionEvalException): def __init__(self) -> None: super().__init__("Expression arithmetic error") diff --git a/qtoggleserver/core/expressions/ports.py b/qtoggleserver/core/expressions/ports.py index 893453dc..fc45e9cf 100644 --- a/qtoggleserver/core/expressions/ports.py +++ b/qtoggleserver/core/expressions/ports.py @@ -5,18 +5,20 @@ from qtoggleserver.core.typing import NullablePortValue from .base import EvalContext, EvalResult, Expression, Role -from .exceptions import DisabledPort, PortValueUnavailable, UnexpectedCharacter, UnknownPortId +from .exceptions import DisabledPort, PortAttrUnavailable, PortValueUnavailable, UnexpectedCharacter, UnknownPortId class PortExpression(Expression, metaclass=abc.ABCMeta): - def __init__(self, port_id: str, prefix: str, role: Role) -> None: + def __init__(self, port_id: str, prefix: str, role: Role, attr_name: str | None = None) -> None: super().__init__(role) self.port_id: str = port_id + self.attr_name: str | None = attr_name self.prefix: str = prefix self._cached_port: core_ports.BasePort | None = None def get_port(self) -> core_ports.BasePort | None: + # TODO: test what happens if a port is removed; do we need this `is_removed` check? port = self._cached_port if port is None or port.is_removed(): port = core_ports.get(self.port_id) @@ -30,23 +32,46 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> E sexpression = stripped.rstrip() prefix = sexpression[0] - port_id = sexpression[1:] - - if port_id: - m = re.search(r"[^a-zA-Z0-9_.-]", port_id) - if m: - p = m.start() - raise UnexpectedCharacter(port_id[p], p + pos + 2) - - if prefix == "$": - return PortValue(port_id, prefix, role) - else: # assuming prefix == '@' - return PortRef(port_id, prefix, role) + sub_sexpression = sexpression[1:] + + if sub_sexpression: + parts = sub_sexpression.split(":", 1) + if len(parts) == 2: # port attribute + port_id, attr_name = parts + m = re.search(r"[^a-zA-Z0-9_.-]", port_id) + if m: + p = m.start() + raise UnexpectedCharacter(port_id[p], p + pos + 2) # TODO: is this the correct position? + m = re.search(r"[^a-zA-Z0-9_-]", attr_name) + if m: + p = m.start() + # TODO: is this the correct position? + raise UnexpectedCharacter(attr_name[p], p + pos + len(port_id) + 3) + + if prefix == "$": + return PortAttr(port_id, prefix, role, attr_name) + else: + raise UnexpectedCharacter(prefix, pos) + else: + port_id = sub_sexpression + m = re.search(r"[^a-zA-Z0-9_.-]", port_id) + if m: + p = m.start() + raise UnexpectedCharacter(sub_sexpression[p], p + pos + 2) + + if prefix == "$": + return PortValue(port_id, prefix, role) + elif prefix == "@": + return PortRef(port_id, prefix, role) + else: + raise UnexpectedCharacter(prefix, pos) else: if prefix == "$": return SelfPortValue(self_port_id, prefix, role) - else: # assuming prefix == '@' + elif prefix == "@": return SelfPortRef(self_port_id, prefix, role) + else: + raise UnexpectedCharacter(prefix, pos) class PortValue(PortExpression): @@ -94,3 +119,20 @@ async def _eval(self, context: EvalContext) -> EvalResult: class SelfPortRef(PortRef): def __str__(self) -> str: return self.prefix + + +class PortAttr(PortExpression): + def __str__(self) -> str: + return f"{self.prefix}{self.port_id}:{self.attr_name}" + + def _get_deps(self) -> set[str]: + return {f"${self.port_id}:"} + + async def _eval(self, context: EvalContext) -> EvalResult: + port = self.get_port() + if not port: + raise UnknownPortId(self.port_id) + + value = context.port_attrs.get(self.port_id, {}).get(self.attr_name) + if value is None: + raise PortAttrUnavailable(self.port_id, self.attr_name or "") From 95358126af387e659994edbde7fd963922dc8278 Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 4 May 2026 11:58:51 +0300 Subject: [PATCH 02/14] core/expressions: Add test cases for new parsing code --- qtoggleserver/core/expressions/devices.py | 4 +- qtoggleserver/core/expressions/ports.py | 7 +-- .../core/expressions/test_parse.py | 48 ++++++++++++++++-- .../core/expressions/test_port.py | 50 ++++++++++++++++++- 4 files changed, 98 insertions(+), 11 deletions(-) diff --git a/qtoggleserver/core/expressions/devices.py b/qtoggleserver/core/expressions/devices.py index 9956482e..0e58ad7c 100644 --- a/qtoggleserver/core/expressions/devices.py +++ b/qtoggleserver/core/expressions/devices.py @@ -27,14 +27,14 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> E if prefix == "#": parts = sub_sexpression.split(":", 1) if len(parts) != 2: - raise MissingAttrPrefix(pos + len(sub_sexpression)) # TODO unit test + raise MissingAttrPrefix(pos + len(sub_sexpression) + 1) device_name, attr_name = parts if device_name: m = re.search(r"[^a-zA-Z0-9_-]", device_name) if m: p = m.start() - raise UnexpectedCharacter(device_name[p], p + pos + 3) # TODO: is this reported position correct? + raise UnexpectedCharacter(device_name[p], p + pos + 3) return SlaveDeviceAttr(device_name, prefix, role, attr_name) else: diff --git a/qtoggleserver/core/expressions/ports.py b/qtoggleserver/core/expressions/ports.py index fc45e9cf..26720e29 100644 --- a/qtoggleserver/core/expressions/ports.py +++ b/qtoggleserver/core/expressions/ports.py @@ -41,12 +41,11 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> E m = re.search(r"[^a-zA-Z0-9_.-]", port_id) if m: p = m.start() - raise UnexpectedCharacter(port_id[p], p + pos + 2) # TODO: is this the correct position? + raise UnexpectedCharacter(port_id[p], p + pos + 2) m = re.search(r"[^a-zA-Z0-9_-]", attr_name) if m: p = m.start() - # TODO: is this the correct position? - raise UnexpectedCharacter(attr_name[p], p + pos + len(port_id) + 3) + raise UnexpectedCharacter(attr_name[p], p + pos + len(port_id) + 2) if prefix == "$": return PortAttr(port_id, prefix, role, attr_name) @@ -136,3 +135,5 @@ async def _eval(self, context: EvalContext) -> EvalResult: value = context.port_attrs.get(self.port_id, {}).get(self.attr_name) if value is None: raise PortAttrUnavailable(self.port_id, self.attr_name or "") + + return value diff --git a/tests/unit/qtoggleserver/core/expressions/test_parse.py b/tests/unit/qtoggleserver/core/expressions/test_parse.py index dd8c0c5a..fd2a2506 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_parse.py +++ b/tests/unit/qtoggleserver/core/expressions/test_parse.py @@ -3,6 +3,7 @@ from qtoggleserver.core.expressions import EvalContext, Role, parse from qtoggleserver.core.expressions.exceptions import ( EmptyExpression, + MissingAttrPrefix, UnbalancedParentheses, UnexpectedCharacter, UnexpectedEnd, @@ -67,9 +68,9 @@ async def test_parse_unexpected_character(): assert exc_info.value.pos == 1 with pytest.raises(UnexpectedCharacter) as exc_info: - parse(None, "ADD#(10, $)", role=Role.VALUE) + parse(None, "ADD?(10, $)", role=Role.VALUE) - assert exc_info.value.c == "#" + assert exc_info.value.c == "?" assert exc_info.value.pos == 4 with pytest.raises(UnexpectedCharacter) as exc_info: @@ -91,9 +92,9 @@ async def test_parse_unexpected_character(): assert exc_info.value.pos == 11 with pytest.raises(UnexpectedCharacter) as exc_info: - parse(None, "ADD(#, 10, $)", role=Role.VALUE) + parse(None, "ADD(?, 10, $)", role=Role.VALUE) - assert exc_info.value.c == "#" + assert exc_info.value.c == "?" assert exc_info.value.pos == 5 with pytest.raises(UnexpectedCharacter) as exc_info: @@ -102,6 +103,45 @@ async def test_parse_unexpected_character(): assert exc_info.value.c == "+" assert exc_info.value.pos == 2 + # `@` prefix is not valid for port-attribute syntax (`$port_id:attr_name`) + with pytest.raises(UnexpectedCharacter) as exc_info: + parse(None, "@nid:attr", role=Role.VALUE) + + assert exc_info.value.c == "@" + assert exc_info.value.pos == 1 + + # invalid character in attr_name part of port-attribute expression + with pytest.raises(UnexpectedCharacter) as exc_info: + parse(None, "$nid:attr*", role=Role.VALUE) + + assert exc_info.value.c == "*" + assert exc_info.value.pos == 10 + + # invalid character in device_name part of device-attribute expression + with pytest.raises(UnexpectedCharacter) as exc_info: + parse(None, "#dev*:attr", role=Role.VALUE) + + assert exc_info.value.c == "*" + + +async def test_parse_missing_attr_prefix(): + # device expression without a colon (no attribute name separator) + with pytest.raises(MissingAttrPrefix) as exc_info: + parse(None, "#device_name", role=Role.VALUE) + + assert exc_info.value.pos == 13 + + with pytest.raises(MissingAttrPrefix) as exc_info: + parse(None, "#", role=Role.VALUE) + + assert exc_info.value.pos == 2 + + # same inside a function argument + with pytest.raises(MissingAttrPrefix) as exc_info: + parse(None, "ADD(#, 10, $)", role=Role.VALUE) + + assert exc_info.value.pos == 6 + async def test_parse_empty(): with pytest.raises(EmptyExpression): diff --git a/tests/unit/qtoggleserver/core/expressions/test_port.py b/tests/unit/qtoggleserver/core/expressions/test_port.py index 08106288..8014eb60 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_port.py +++ b/tests/unit/qtoggleserver/core/expressions/test_port.py @@ -1,8 +1,13 @@ import pytest from qtoggleserver.core.expressions import Role, parse -from qtoggleserver.core.expressions.exceptions import DisabledPort, PortValueUnavailable, UnknownPortId -from qtoggleserver.core.expressions.ports import PortRef, PortValue, SelfPortRef, SelfPortValue +from qtoggleserver.core.expressions.exceptions import ( + DisabledPort, + PortAttrUnavailable, + PortValueUnavailable, + UnknownPortId, +) +from qtoggleserver.core.expressions.ports import PortAttr, PortRef, PortValue, SelfPortRef, SelfPortValue class TestPortValue: @@ -121,3 +126,44 @@ async def test_cache_invalidated_on_port_removal(self, mock_num_port1): await mock_num_port1.remove(persisted_data=False) assert mock_num_port1.is_removed() assert e.get_port() is None + + +class TestPortAttr: + def test_parse(self, mock_num_port1): + """Should parse `$port_id:attr_name` into a PortAttr with the correct port_id and attr_name.""" + + e = parse("nid1", "$nid1:enabled", role=Role.VALUE) + assert isinstance(e, PortAttr) + assert e.port_id == "nid1" + assert e.attr_name == "enabled" + + def test_deps(self): + """Should return a dep set with `$port_id:` format.""" + + e = PortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="enabled") + assert e._get_deps() == {"$nid1:"} + + async def test_eval(self, mock_num_port1, dummy_eval_context): + """Should return the attribute value from the supplied context.""" + + e = PortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="enabled") + dummy_eval_context.port_attrs["nid1"] = {"enabled": True} + assert await e._eval(dummy_eval_context) is True + + async def test_eval_unknown(self, dummy_eval_context): + """Should raise UnknownPortId when the referenced port doesn't exist.""" + + e = PortAttr("inexistent", prefix="$", role=Role.VALUE, attr_name="enabled") + with pytest.raises(UnknownPortId) as exc_info: + await e._eval(dummy_eval_context) + assert exc_info.value.port_id == "inexistent" + + async def test_eval_unavailable(self, mock_num_port1, dummy_eval_context): + """Should raise PortAttrUnavailable when the attribute is not present in the context.""" + + e = PortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="nonexistent_attr") + dummy_eval_context.port_attrs["nid1"] = {} + with pytest.raises(PortAttrUnavailable) as exc_info: + await e._eval(dummy_eval_context) + assert exc_info.value.port_id == "nid1" + assert exc_info.value.attr_name == "nonexistent_attr" From 250c90eb764cddf7e62cb1e0028bad1818d2d65f Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 4 May 2026 12:01:59 +0300 Subject: [PATCH 03/14] core/expressions: Move functions-specific test cases to functions package --- tests/unit/qtoggleserver/core/expressions/functions/__init__.py | 0 .../core/expressions/{ => functions}/test_aggregation.py | 0 .../core/expressions/{ => functions}/test_arithmetic.py | 0 .../core/expressions/{ => functions}/test_bitwise.py | 0 .../core/expressions/{ => functions}/test_branching.py | 0 .../core/expressions/{ => functions}/test_comparison.py | 0 .../qtoggleserver/core/expressions/{ => functions}/test_date.py | 0 .../core/expressions/{ => functions}/test_function.py | 0 .../qtoggleserver/core/expressions/{ => functions}/test_logic.py | 0 .../core/expressions/{ => functions}/test_rounding.py | 0 .../qtoggleserver/core/expressions/{ => functions}/test_sign.py | 0 .../qtoggleserver/core/expressions/{ => functions}/test_time.py | 0 .../core/expressions/{ => functions}/test_timeprocessing.py | 0 .../core/expressions/{ => functions}/test_various.py | 0 14 files changed, 0 insertions(+), 0 deletions(-) create mode 100644 tests/unit/qtoggleserver/core/expressions/functions/__init__.py rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_aggregation.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_arithmetic.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_bitwise.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_branching.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_comparison.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_date.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_function.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_logic.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_rounding.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_sign.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_time.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_timeprocessing.py (100%) rename tests/unit/qtoggleserver/core/expressions/{ => functions}/test_various.py (100%) diff --git a/tests/unit/qtoggleserver/core/expressions/functions/__init__.py b/tests/unit/qtoggleserver/core/expressions/functions/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/tests/unit/qtoggleserver/core/expressions/test_aggregation.py b/tests/unit/qtoggleserver/core/expressions/functions/test_aggregation.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_aggregation.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_aggregation.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_arithmetic.py b/tests/unit/qtoggleserver/core/expressions/functions/test_arithmetic.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_arithmetic.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_arithmetic.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_bitwise.py b/tests/unit/qtoggleserver/core/expressions/functions/test_bitwise.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_bitwise.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_bitwise.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_branching.py b/tests/unit/qtoggleserver/core/expressions/functions/test_branching.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_branching.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_branching.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_comparison.py b/tests/unit/qtoggleserver/core/expressions/functions/test_comparison.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_comparison.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_comparison.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_date.py b/tests/unit/qtoggleserver/core/expressions/functions/test_date.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_date.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_date.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_function.py b/tests/unit/qtoggleserver/core/expressions/functions/test_function.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_function.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_function.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_logic.py b/tests/unit/qtoggleserver/core/expressions/functions/test_logic.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_logic.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_logic.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_rounding.py b/tests/unit/qtoggleserver/core/expressions/functions/test_rounding.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_rounding.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_rounding.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_sign.py b/tests/unit/qtoggleserver/core/expressions/functions/test_sign.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_sign.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_sign.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_time.py b/tests/unit/qtoggleserver/core/expressions/functions/test_time.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_time.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_time.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_timeprocessing.py b/tests/unit/qtoggleserver/core/expressions/functions/test_timeprocessing.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_timeprocessing.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_timeprocessing.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_various.py b/tests/unit/qtoggleserver/core/expressions/functions/test_various.py similarity index 100% rename from tests/unit/qtoggleserver/core/expressions/test_various.py rename to tests/unit/qtoggleserver/core/expressions/functions/test_various.py From 87277e8f7f141679d18983ca7e14b99a20bd5c01 Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 4 May 2026 12:07:20 +0300 Subject: [PATCH 04/14] core/expressions: Add device expressions test cases --- .../core/expressions/test_device.py | 81 +++++++++++++++++++ 1 file changed, 81 insertions(+) create mode 100644 tests/unit/qtoggleserver/core/expressions/test_device.py diff --git a/tests/unit/qtoggleserver/core/expressions/test_device.py b/tests/unit/qtoggleserver/core/expressions/test_device.py new file mode 100644 index 00000000..2cf10a3f --- /dev/null +++ b/tests/unit/qtoggleserver/core/expressions/test_device.py @@ -0,0 +1,81 @@ +import pytest + +from qtoggleserver.core.expressions import Role, parse +from qtoggleserver.core.expressions.devices import MainDeviceAttr, SlaveDeviceAttr +from qtoggleserver.core.expressions.exceptions import DeviceAttrUnavailable + + +class TestMainDeviceAttr: + def test_parse(self): + """Should parse `#:attr_name` into a MainDeviceAttr with device_name=None.""" + + e = parse(None, "#:display_name", role=Role.VALUE) + assert isinstance(e, MainDeviceAttr) + assert e.device_name is None + assert e.attr_name == "display_name" + + def test_deps(self): + """Should return `{'#:'}` as a dependency.""" + + e = MainDeviceAttr(prefix="#", role=Role.VALUE, attr_name="display_name") + assert e._get_deps() == {"#:"} + + async def test_eval(self, dummy_eval_context): + """Should return the attribute value from the supplied context.""" + + e = MainDeviceAttr(prefix="#", role=Role.VALUE, attr_name="main_test_attr") + dummy_eval_context.device_attrs["main_test_attr"] = "My Device" + assert await e._eval(dummy_eval_context) == "My Device" + + async def test_eval_unavailable(self, dummy_eval_context): + """Should raise DeviceAttrUnavailable when the attribute is not present in the context.""" + + e = MainDeviceAttr(prefix="#", role=Role.VALUE, attr_name="nonexistent_main_attr") + with pytest.raises(DeviceAttrUnavailable) as exc_info: + await e._eval(dummy_eval_context) + assert exc_info.value.device_name == "" + assert exc_info.value.attr_name == "nonexistent_main_attr" + + def test_str(self): + """Should serialise back to the original expression syntax.""" + + e = MainDeviceAttr(prefix="#", role=Role.VALUE, attr_name="display_name") + assert str(e) == "#:display_name" + + +class TestSlaveDeviceAttr: + def test_parse(self): + """Should parse `#slave_name:attr_name` into a SlaveDeviceAttr.""" + + e = parse(None, "#slave1:display_name", role=Role.VALUE) + assert isinstance(e, SlaveDeviceAttr) + assert e.device_name == "slave1" + assert e.attr_name == "display_name" + + def test_deps(self): + """Should return `{'#slave_name:'}` as a dependency.""" + + e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="display_name") + assert e._get_deps() == {"#slave1:"} + + async def test_eval(self, dummy_eval_context): + """Should return the attribute value using a `slave_name.attr_name` key in the context.""" + + e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="display_name") + dummy_eval_context.device_attrs["slave1.display_name"] = "Slave Device" + assert await e._eval(dummy_eval_context) == "Slave Device" + + async def test_eval_unavailable(self, dummy_eval_context): + """Should raise DeviceAttrUnavailable when the attribute is not present in the context.""" + + e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="nonexistent_slave_attr") + with pytest.raises(DeviceAttrUnavailable) as exc_info: + await e._eval(dummy_eval_context) + assert exc_info.value.device_name == "slave1" + assert exc_info.value.attr_name == "nonexistent_slave_attr" + + def test_str(self): + """Should serialise back to the original expression syntax.""" + + e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="display_name") + assert str(e) == "#slave1:display_name" From c4d9cd2a0bacd383d8350e2f9d5f2626023c074c Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 4 May 2026 14:33:13 +0300 Subject: [PATCH 05/14] core/expressions: Add SelfPortAttr, enforce transform limitations --- qtoggleserver/core/expressions/devices.py | 9 ++- qtoggleserver/core/expressions/exceptions.py | 22 +++--- qtoggleserver/core/expressions/ports.py | 30 ++++++++- qtoggleserver/core/ports.py | 10 --- qtoggleserver/utils/expressions.py | 2 +- .../core/expressions/test_device.py | 2 +- .../core/expressions/test_parse.py | 7 ++ .../core/expressions/test_port.py | 49 +++++++++++++- .../core/expressions/test_transform.py | 67 +++++++++++++++++-- .../qtoggleserver/utils/test_expressions.py | 8 +-- 10 files changed, 168 insertions(+), 38 deletions(-) diff --git a/qtoggleserver/core/expressions/devices.py b/qtoggleserver/core/expressions/devices.py index 0e58ad7c..ec88d87e 100644 --- a/qtoggleserver/core/expressions/devices.py +++ b/qtoggleserver/core/expressions/devices.py @@ -4,7 +4,10 @@ import re from .base import EvalContext, EvalResult, Expression, Role -from .exceptions import DeviceAttrUnavailable, MissingAttrPrefix, UnexpectedCharacter +from .exceptions import DeviceAttrUnavailable, MissingAttrPrefix, TransformNotSupported, UnexpectedCharacter + + +_TRANSFORM_ROLES = (Role.TRANSFORM_READ, Role.TRANSFORM_WRITE) class DeviceExpression(Expression, metaclass=abc.ABCMeta): @@ -25,6 +28,8 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> E sub_sexpression = sexpression[1:] if prefix == "#": + if role in _TRANSFORM_ROLES: + raise TransformNotSupported(sexpression, pos) parts = sub_sexpression.split(":", 1) if len(parts) != 2: raise MissingAttrPrefix(pos + len(sub_sexpression) + 1) @@ -51,7 +56,7 @@ def _get_deps(self) -> set[str]: return {f"#{self.device_name or ''}:"} async def _eval(self, context: EvalContext) -> EvalResult: - key = f"{self.device_name}.{self.attr_name or ''}" if self.device_name else self.attr_name + key = f"{self.device_name}:{self.attr_name or ''}" if self.device_name else self.attr_name value = context.device_attrs.get(key) if value is None: raise DeviceAttrUnavailable(self.device_name or "", self.attr_name or "") diff --git a/qtoggleserver/core/expressions/exceptions.py b/qtoggleserver/core/expressions/exceptions.py index b006ee74..d7a09ae2 100644 --- a/qtoggleserver/core/expressions/exceptions.py +++ b/qtoggleserver/core/expressions/exceptions.py @@ -63,17 +63,6 @@ def to_json(self) -> GenericJSONDict: return {"reason": "unexpected-end"} -class NonSelfDependency(ExpressionParseError): - def __init__(self, port_id: str, pos: int) -> None: - self.port_id = port_id - self.pos = pos - - super().__init__(f'Non-self dependency on port "{port_id}"') - - def to_json(self) -> GenericJSONDict: - return {"reason": "non-self-dependency", "token": self.port_id, "pos": self.pos} - - class UnexpectedCharacter(ExpressionParseError): def __init__(self, c: str, pos: int) -> None: self.c = c @@ -103,6 +92,17 @@ def to_json(self) -> GenericJSONDict: return {"reason": "missing-attr-prefix", "pos": self.pos} +class TransformNotSupported(ExpressionParseError): + def __init__(self, token: str, pos: int) -> None: + self.token: str = token + self.pos: int = pos + + super().__init__(f'Expression "{token}" is not supported in transform expressions') + + def to_json(self) -> GenericJSONDict: + return {"reason": "transform-not-supported", "token": self.token, "pos": self.pos} + + class ExpressionEvalException(ExpressionException): pass diff --git a/qtoggleserver/core/expressions/ports.py b/qtoggleserver/core/expressions/ports.py index 26720e29..2d0c3e00 100644 --- a/qtoggleserver/core/expressions/ports.py +++ b/qtoggleserver/core/expressions/ports.py @@ -5,7 +5,17 @@ from qtoggleserver.core.typing import NullablePortValue from .base import EvalContext, EvalResult, Expression, Role -from .exceptions import DisabledPort, PortAttrUnavailable, PortValueUnavailable, UnexpectedCharacter, UnknownPortId +from .exceptions import ( + DisabledPort, + PortAttrUnavailable, + PortValueUnavailable, + TransformNotSupported, + UnexpectedCharacter, + UnknownPortId, +) + + +_TRANSFORM_ROLES = (Role.TRANSFORM_READ, Role.TRANSFORM_WRITE) class PortExpression(Expression, metaclass=abc.ABCMeta): @@ -48,7 +58,12 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> E raise UnexpectedCharacter(attr_name[p], p + pos + len(port_id) + 2) if prefix == "$": - return PortAttr(port_id, prefix, role, attr_name) + if role in _TRANSFORM_ROLES: + raise TransformNotSupported(sexpression, pos) + if port_id: + return PortAttr(port_id, prefix, role, attr_name) + else: + return SelfPortAttr(self_port_id, prefix, role, attr_name) else: raise UnexpectedCharacter(prefix, pos) else: @@ -59,8 +74,12 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> E raise UnexpectedCharacter(sub_sexpression[p], p + pos + 2) if prefix == "$": + if role in _TRANSFORM_ROLES and port_id != self_port_id: + raise TransformNotSupported(sexpression, pos=pos) return PortValue(port_id, prefix, role) elif prefix == "@": + if role in _TRANSFORM_ROLES: + raise TransformNotSupported(sexpression, pos) return PortRef(port_id, prefix, role) else: raise UnexpectedCharacter(prefix, pos) @@ -68,6 +87,8 @@ def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> E if prefix == "$": return SelfPortValue(self_port_id, prefix, role) elif prefix == "@": + if role in _TRANSFORM_ROLES: + raise TransformNotSupported(sexpression, pos) return SelfPortRef(self_port_id, prefix, role) else: raise UnexpectedCharacter(prefix, pos) @@ -137,3 +158,8 @@ async def _eval(self, context: EvalContext) -> EvalResult: raise PortAttrUnavailable(self.port_id, self.attr_name or "") return value + + +class SelfPortAttr(PortAttr): + def __str__(self) -> str: + return f"{self.prefix}:{self.attr_name}" diff --git a/qtoggleserver/core/ports.py b/qtoggleserver/core/ports.py index a230770a..3faca488 100644 --- a/qtoggleserver/core/ports.py +++ b/qtoggleserver/core/ports.py @@ -599,11 +599,6 @@ async def attr_set_transform_read(self, stransform_read: str) -> None: self.get_id(), stransform_read, role=core_expressions.Role.TRANSFORM_READ ) - deps = transform_read.get_deps() - for dep in deps: - if dep.startswith("$") and len(dep) > 1 and dep[1:] != self._id: - raise expressions_exceptions.NonSelfDependency(port_id=dep[1:], pos=stransform_read.index(dep)) - self.debug('setting read transform "%s"', transform_read) self._transform_read = transform_read except expressions_exceptions.ExpressionParseError as e: @@ -635,11 +630,6 @@ async def attr_set_transform_write(self, stransform_write: str) -> None: self.get_id(), stransform_write, role=core_expressions.Role.TRANSFORM_WRITE ) - deps = transform_write.get_deps() - for dep in deps: - if dep.startswith("$") and len(dep) > 1 and dep[1:] != self._id: - raise expressions_exceptions.NonSelfDependency(port_id=dep[1:], pos=stransform_write.index(dep)) - self.debug('setting write transform "%s"', transform_write) self._transform_write = transform_write except expressions_exceptions.ExpressionParseError as e: diff --git a/qtoggleserver/utils/expressions.py b/qtoggleserver/utils/expressions.py index f8d65e74..4f50a6c1 100644 --- a/qtoggleserver/utils/expressions.py +++ b/qtoggleserver/utils/expressions.py @@ -66,7 +66,7 @@ async def build_context(now_ms: int) -> EvalContext: slave_name = slave.get_name() slave_attrs = slave.get_cached_attrs() device_attrs.update( - {f"{slave_name}.{attr_name}": attr_value for attr_name, attr_value in slave_attrs.items()} + {f"{slave_name}:{attr_name}": attr_value for attr_name, attr_value in slave_attrs.items()} ) return EvalContext(port_values, port_attrs, device_attrs, now_ms) diff --git a/tests/unit/qtoggleserver/core/expressions/test_device.py b/tests/unit/qtoggleserver/core/expressions/test_device.py index 2cf10a3f..ef6e2fcd 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_device.py +++ b/tests/unit/qtoggleserver/core/expressions/test_device.py @@ -62,7 +62,7 @@ async def test_eval(self, dummy_eval_context): """Should return the attribute value using a `slave_name.attr_name` key in the context.""" e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="display_name") - dummy_eval_context.device_attrs["slave1.display_name"] = "Slave Device" + dummy_eval_context.device_attrs["slave1:display_name"] = "Slave Device" assert await e._eval(dummy_eval_context) == "Slave Device" async def test_eval_unavailable(self, dummy_eval_context): diff --git a/tests/unit/qtoggleserver/core/expressions/test_parse.py b/tests/unit/qtoggleserver/core/expressions/test_parse.py index fd2a2506..37f1c5d0 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_parse.py +++ b/tests/unit/qtoggleserver/core/expressions/test_parse.py @@ -117,6 +117,13 @@ async def test_parse_unexpected_character(): assert exc_info.value.c == "*" assert exc_info.value.pos == 10 + # invalid character in attr_name part of self-port-attribute expression + with pytest.raises(UnexpectedCharacter) as exc_info: + parse(None, "$:attr*", role=Role.VALUE) + + assert exc_info.value.c == "*" + assert exc_info.value.pos == 7 + # invalid character in device_name part of device-attribute expression with pytest.raises(UnexpectedCharacter) as exc_info: parse(None, "#dev*:attr", role=Role.VALUE) diff --git a/tests/unit/qtoggleserver/core/expressions/test_port.py b/tests/unit/qtoggleserver/core/expressions/test_port.py index 8014eb60..868be169 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_port.py +++ b/tests/unit/qtoggleserver/core/expressions/test_port.py @@ -7,7 +7,7 @@ PortValueUnavailable, UnknownPortId, ) -from qtoggleserver.core.expressions.ports import PortAttr, PortRef, PortValue, SelfPortRef, SelfPortValue +from qtoggleserver.core.expressions.ports import PortAttr, PortRef, PortValue, SelfPortAttr, SelfPortRef, SelfPortValue class TestPortValue: @@ -167,3 +167,50 @@ async def test_eval_unavailable(self, mock_num_port1, dummy_eval_context): await e._eval(dummy_eval_context) assert exc_info.value.port_id == "nid1" assert exc_info.value.attr_name == "nonexistent_attr" + + +class TestSelfPortAttr: + def test_parse(self, mock_num_port1): + """Should parse `$:attr_name` into a SelfPortAttr with port_id set to self_port_id.""" + + e = parse("nid1", "$:enabled", role=Role.VALUE) + assert isinstance(e, SelfPortAttr) + assert e.port_id == "nid1" + assert e.attr_name == "enabled" + + def test_str(self): + """Should render as `$:attr_name`.""" + + e = SelfPortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="enabled") + assert str(e) == "$:enabled" + + def test_deps(self): + """Should return a dep set referencing the self port's attributes.""" + + e = SelfPortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="enabled") + assert e._get_deps() == {"$nid1:"} + + async def test_eval(self, mock_num_port1, dummy_eval_context): + """Should return the attribute value of the self port from the supplied context.""" + + e = SelfPortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="enabled") + dummy_eval_context.port_attrs["nid1"] = {"enabled": True} + assert await e._eval(dummy_eval_context) is True + + async def test_eval_unknown(self, dummy_eval_context): + """Should raise UnknownPortId when the self port doesn't exist.""" + + e = SelfPortAttr("inexistent", prefix="$", role=Role.VALUE, attr_name="enabled") + with pytest.raises(UnknownPortId) as exc_info: + await e._eval(dummy_eval_context) + assert exc_info.value.port_id == "inexistent" + + async def test_eval_unavailable(self, mock_num_port1, dummy_eval_context): + """Should raise PortAttrUnavailable when the attribute is not present in the context.""" + + e = SelfPortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="nonexistent_attr") + dummy_eval_context.port_attrs["nid1"] = {} + with pytest.raises(PortAttrUnavailable) as exc_info: + await e._eval(dummy_eval_context) + assert exc_info.value.port_id == "nid1" + assert exc_info.value.attr_name == "nonexistent_attr" diff --git a/tests/unit/qtoggleserver/core/expressions/test_transform.py b/tests/unit/qtoggleserver/core/expressions/test_transform.py index 704c3d75..034466ac 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_transform.py +++ b/tests/unit/qtoggleserver/core/expressions/test_transform.py @@ -10,12 +10,51 @@ async def test(self, mock_num_port1): mock_num_port1.set_next_value(5) assert await mock_num_port1.read_transformed_value() == 50 - async def test_non_self_dependency(self, mock_num_port1): - """Should (indirectly) raise `NonSelfDependency` when expression depends on another port.""" + async def test_transform_not_supported(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression depends on another port.""" with pytest.raises(InvalidAttributeValue) as exc_info: await mock_num_port1.set_attr("transform_read", "MUL($another_port, 10)") - assert exc_info.value.details == {"reason": "non-self-dependency", "token": "another_port", "pos": 4} + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "$another_port", "pos": 5} + + async def test_port_reference_forbidden(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression contains a port reference.""" + + with pytest.raises(InvalidAttributeValue) as exc_info: + await mock_num_port1.set_attr("transform_read", "@another_port") + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "@another_port", "pos": 1} + + async def test_self_port_reference_forbidden(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression contains a self port reference.""" + + with pytest.raises(InvalidAttributeValue) as exc_info: + await mock_num_port1.set_attr("transform_read", "@") + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "@", "pos": 1} + + async def test_port_attribute_forbidden(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression contains a port attribute.""" + + with pytest.raises(InvalidAttributeValue) as exc_info: + await mock_num_port1.set_attr("transform_read", "$another_port:enabled") + assert exc_info.value.details == { + "reason": "transform-not-supported", + "token": "$another_port:enabled", + "pos": 1, + } + + async def test_self_port_attribute_forbidden(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression contains a self port attribute.""" + + with pytest.raises(InvalidAttributeValue) as exc_info: + await mock_num_port1.set_attr("transform_read", "$:enabled") + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "$:enabled", "pos": 1} + + async def test_device_attribute_forbidden(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression contains a device attribute.""" + + with pytest.raises(InvalidAttributeValue) as exc_info: + await mock_num_port1.set_attr("transform_read", "#:name") + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "#:name", "pos": 1} class TestPortTransformWrite: @@ -30,10 +69,26 @@ async def test(self, mock_num_port1, mocker): await mock_num_port1.transform_and_write_value(6) mock_num_port1.write_value.assert_called_once_with(60) - async def test_non_self_dependency(self, mock_num_port1): - """Should (indirectly) raise `NonSelfDependency` when expression depends on another port.""" + async def test_transform_not_supported(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression depends on another port.""" mock_num_port1.set_writable(True) with pytest.raises(InvalidAttributeValue) as exc_info: await mock_num_port1.set_attr("transform_write", "MUL($another_port, 10)") - assert exc_info.value.details == {"reason": "non-self-dependency", "token": "another_port", "pos": 4} + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "$another_port", "pos": 5} + + async def test_port_reference_forbidden(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression contains a port reference.""" + + mock_num_port1.set_writable(True) + with pytest.raises(InvalidAttributeValue) as exc_info: + await mock_num_port1.set_attr("transform_write", "@another_port") + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "@another_port", "pos": 1} + + async def test_device_attribute_forbidden(self, mock_num_port1): + """Should raise `TransformNotSupported` when expression contains a device attribute.""" + + mock_num_port1.set_writable(True) + with pytest.raises(InvalidAttributeValue) as exc_info: + await mock_num_port1.set_attr("transform_write", "#:name") + assert exc_info.value.details == {"reason": "transform-not-supported", "token": "#:name", "pos": 1} diff --git a/tests/unit/qtoggleserver/utils/test_expressions.py b/tests/unit/qtoggleserver/utils/test_expressions.py index 4ef923d0..1875f7e7 100644 --- a/tests/unit/qtoggleserver/utils/test_expressions.py +++ b/tests/unit/qtoggleserver/utils/test_expressions.py @@ -191,8 +191,8 @@ async def test_with_slave_attrs(self, mock_num_port1, mocker): assert context.device_attrs == { "device_attr": "device_val", - "slave1.attr1": "val1", - "slave1.attr2": "val2", + "slave1:attr1": "val1", + "slave1:attr2": "val2", } assert context.now_ms == 3000 @@ -228,8 +228,8 @@ async def test_with_multiple_slaves(self, mock_num_port1, mocker): assert context.device_attrs == { "device_attr": "device_val", - "slave1.attr1": "val1", - "slave2.attr2": "val2", + "slave1:attr1": "val1", + "slave2:attr2": "val2", } assert context.now_ms == 4000 From b7c8615991432818ab2c140bdb15cf112890f766 Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 4 May 2026 21:20:11 +0300 Subject: [PATCH 06/14] frontend: Update invalid expression reasons --- qtoggleserver/frontend/js/api/constants.js | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/qtoggleserver/frontend/js/api/constants.js b/qtoggleserver/frontend/js/api/constants.js index b91adbf8..f686882c 100644 --- a/qtoggleserver/frontend/js/api/constants.js +++ b/qtoggleserver/frontend/js/api/constants.js @@ -260,14 +260,6 @@ export const INVALID_EXPRESSION_REASONS = [ reason: 'unexpected-end', pretty: gettext('Expression is unterminated.') }, - { - reason: 'circular-dependency', - pretty: gettext('Expression creates a dependency loop.') - }, - { - reason: 'external-dependency', - pretty: gettext('Expression must not depend on other ports.') - }, { reason: 'unexpected-character', pretty: gettext('Unexpected character "%(token)s" at position %(pos)d.') @@ -279,6 +271,14 @@ export const INVALID_EXPRESSION_REASONS = [ { reason: 'too-long', pretty: gettext('Expression is too long.') + }, + { + reason: 'missing-attr-prefix', + pretty: gettext('Missing attribute prefix at position %(pos)d.') + }, + { + reason: 'transform-not-supported', + pretty: gettext('Expression "%(token)s" is not supported in transform expressions.') } ] From 61a7adce710d2d9563229bb27920734a82377a2d Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Sun, 17 May 2026 22:06:22 +0300 Subject: [PATCH 07/14] Properly coerce string attribute expressions --- qtoggleserver/core/expressions/devices.py | 2 ++ qtoggleserver/core/expressions/ports.py | 2 ++ .../core/expressions/test_device.py | 32 +++++++++++++++++-- .../core/expressions/test_port.py | 28 ++++++++++++++++ 4 files changed, 62 insertions(+), 2 deletions(-) diff --git a/qtoggleserver/core/expressions/devices.py b/qtoggleserver/core/expressions/devices.py index ec88d87e..cc040005 100644 --- a/qtoggleserver/core/expressions/devices.py +++ b/qtoggleserver/core/expressions/devices.py @@ -60,6 +60,8 @@ async def _eval(self, context: EvalContext) -> EvalResult: value = context.device_attrs.get(key) if value is None: raise DeviceAttrUnavailable(self.device_name or "", self.attr_name or "") + if not isinstance(value, (int, float)): # this includes `bool` + value = int(bool(value)) return value diff --git a/qtoggleserver/core/expressions/ports.py b/qtoggleserver/core/expressions/ports.py index 2d0c3e00..9e43a598 100644 --- a/qtoggleserver/core/expressions/ports.py +++ b/qtoggleserver/core/expressions/ports.py @@ -156,6 +156,8 @@ async def _eval(self, context: EvalContext) -> EvalResult: value = context.port_attrs.get(self.port_id, {}).get(self.attr_name) if value is None: raise PortAttrUnavailable(self.port_id, self.attr_name or "") + if not isinstance(value, (int, float)): # this includes `bool` + value = int(bool(value)) return value diff --git a/tests/unit/qtoggleserver/core/expressions/test_device.py b/tests/unit/qtoggleserver/core/expressions/test_device.py index ef6e2fcd..7706e20c 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_device.py +++ b/tests/unit/qtoggleserver/core/expressions/test_device.py @@ -23,9 +23,23 @@ def test_deps(self): async def test_eval(self, dummy_eval_context): """Should return the attribute value from the supplied context.""" + e = MainDeviceAttr(prefix="#", role=Role.VALUE, attr_name="main_test_attr") + dummy_eval_context.device_attrs["main_test_attr"] = 42 + assert await e._eval(dummy_eval_context) == 42 + + async def test_eval_string_nonempty(self, dummy_eval_context): + """Should return 1 when the attribute is a non-empty string.""" + e = MainDeviceAttr(prefix="#", role=Role.VALUE, attr_name="main_test_attr") dummy_eval_context.device_attrs["main_test_attr"] = "My Device" - assert await e._eval(dummy_eval_context) == "My Device" + assert await e._eval(dummy_eval_context) == 1 + + async def test_eval_string_empty(self, dummy_eval_context): + """Should return 0 when the attribute is an empty string.""" + + e = MainDeviceAttr(prefix="#", role=Role.VALUE, attr_name="main_test_attr") + dummy_eval_context.device_attrs["main_test_attr"] = "" + assert await e._eval(dummy_eval_context) == 0 async def test_eval_unavailable(self, dummy_eval_context): """Should raise DeviceAttrUnavailable when the attribute is not present in the context.""" @@ -61,9 +75,23 @@ def test_deps(self): async def test_eval(self, dummy_eval_context): """Should return the attribute value using a `slave_name.attr_name` key in the context.""" + e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="display_name") + dummy_eval_context.device_attrs["slave1:display_name"] = 42 + assert await e._eval(dummy_eval_context) == 42 + + async def test_eval_string_nonempty(self, dummy_eval_context): + """Should return 1 when the attribute is a non-empty string.""" + e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="display_name") dummy_eval_context.device_attrs["slave1:display_name"] = "Slave Device" - assert await e._eval(dummy_eval_context) == "Slave Device" + assert await e._eval(dummy_eval_context) == 1 + + async def test_eval_string_empty(self, dummy_eval_context): + """Should return 0 when the attribute is an empty string.""" + + e = SlaveDeviceAttr("slave1", prefix="#", role=Role.VALUE, attr_name="display_name") + dummy_eval_context.device_attrs["slave1:display_name"] = "" + assert await e._eval(dummy_eval_context) == 0 async def test_eval_unavailable(self, dummy_eval_context): """Should raise DeviceAttrUnavailable when the attribute is not present in the context.""" diff --git a/tests/unit/qtoggleserver/core/expressions/test_port.py b/tests/unit/qtoggleserver/core/expressions/test_port.py index 868be169..fe501657 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_port.py +++ b/tests/unit/qtoggleserver/core/expressions/test_port.py @@ -150,6 +150,20 @@ async def test_eval(self, mock_num_port1, dummy_eval_context): dummy_eval_context.port_attrs["nid1"] = {"enabled": True} assert await e._eval(dummy_eval_context) is True + async def test_eval_string_nonempty(self, mock_num_port1, dummy_eval_context): + """Should return 1 when the attribute is a non-empty string.""" + + e = PortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="tag") + dummy_eval_context.port_attrs["nid1"] = {"tag": "sensor"} + assert await e._eval(dummy_eval_context) == 1 + + async def test_eval_string_empty(self, mock_num_port1, dummy_eval_context): + """Should return 0 when the attribute is an empty string.""" + + e = PortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="tag") + dummy_eval_context.port_attrs["nid1"] = {"tag": ""} + assert await e._eval(dummy_eval_context) == 0 + async def test_eval_unknown(self, dummy_eval_context): """Should raise UnknownPortId when the referenced port doesn't exist.""" @@ -197,6 +211,20 @@ async def test_eval(self, mock_num_port1, dummy_eval_context): dummy_eval_context.port_attrs["nid1"] = {"enabled": True} assert await e._eval(dummy_eval_context) is True + async def test_eval_string_nonempty(self, mock_num_port1, dummy_eval_context): + """Should return 1 when the attribute is a non-empty string.""" + + e = SelfPortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="tag") + dummy_eval_context.port_attrs["nid1"] = {"tag": "sensor"} + assert await e._eval(dummy_eval_context) == 1 + + async def test_eval_string_empty(self, mock_num_port1, dummy_eval_context): + """Should return 0 when the attribute is an empty string.""" + + e = SelfPortAttr("nid1", prefix="$", role=Role.VALUE, attr_name="tag") + dummy_eval_context.port_attrs["nid1"] = {"tag": ""} + assert await e._eval(dummy_eval_context) == 0 + async def test_eval_unknown(self, dummy_eval_context): """Should raise UnknownPortId when the self port doesn't exist.""" From 2e384bb5cd5d2dce180a85e533bda7ce2dc3475f Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Sun, 17 May 2026 22:48:15 +0300 Subject: [PATCH 08/14] peripherals: Fix set_online bug --- qtoggleserver/peripherals/peripheral.py | 4 +- .../peripherals/test_peripheral.py | 66 +++++++++++++++++++ 2 files changed, 68 insertions(+), 2 deletions(-) create mode 100644 tests/unit/qtoggleserver/peripherals/test_peripheral.py diff --git a/qtoggleserver/peripherals/peripheral.py b/qtoggleserver/peripherals/peripheral.py index 08352a7c..9f3f7571 100644 --- a/qtoggleserver/peripherals/peripheral.py +++ b/qtoggleserver/peripherals/peripheral.py @@ -168,16 +168,16 @@ def is_online(self) -> bool: return self._enabled and self._online def set_online(self, online: bool) -> None: - self._online = online - if online and not self._online: self.debug("is online") + self._online = online try: self.handle_online() except Exception: self.error("handle_online failed", exc_info=True) elif not online and self._online: self.debug("is offline") + self._online = online try: self.handle_offline() except Exception: diff --git a/tests/unit/qtoggleserver/peripherals/test_peripheral.py b/tests/unit/qtoggleserver/peripherals/test_peripheral.py new file mode 100644 index 00000000..2b98aec4 --- /dev/null +++ b/tests/unit/qtoggleserver/peripherals/test_peripheral.py @@ -0,0 +1,66 @@ +from tests.unit.qtoggleserver.mock.peripherals import MockPeripheral + + +class TestSetOnline: + def make_peripheral(self, mocker) -> MockPeripheral: + p = MockPeripheral(name="test", dummy_param="v") + mocker.patch.object(p, "handle_online") + mocker.patch.object(p, "handle_offline") + return p + + def test_handle_online_called_when_transitioning_to_online(self, mocker): + """Should call handle_online() exactly once when transitioning from offline to online.""" + + p = self.make_peripheral(mocker) + assert not p._online + + p.set_online(True) + + p.handle_online.assert_called_once() + p.handle_offline.assert_not_called() + + def test_handle_offline_called_when_transitioning_to_offline(self, mocker): + """Should call handle_offline() exactly once when transitioning from online to offline.""" + + p = self.make_peripheral(mocker) + p._online = True + + p.set_online(False) + + p.handle_offline.assert_called_once() + p.handle_online.assert_not_called() + + def test_handle_online_not_called_when_already_online(self, mocker): + """Should not call handle_online() when the peripheral is already online.""" + + p = self.make_peripheral(mocker) + p._online = True + + p.set_online(True) + + p.handle_online.assert_not_called() + + def test_handle_offline_not_called_when_already_offline(self, mocker): + """Should not call handle_offline() when the peripheral is already offline.""" + + p = self.make_peripheral(mocker) + assert not p._online + + p.set_online(False) + + p.handle_offline.assert_not_called() + + def test_online_state_updated_when_going_online(self, mocker): + """Should update _online to True after set_online(True).""" + + p = self.make_peripheral(mocker) + p.set_online(True) + assert p._online is True + + def test_online_state_updated_when_going_offline(self, mocker): + """Should update _online to False after set_online(False).""" + + p = self.make_peripheral(mocker) + p._online = True + p.set_online(False) + assert p._online is False From 4d26df60aae6d34e2070ba4db84fbb5f76caf4df Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Sun, 17 May 2026 22:50:23 +0300 Subject: [PATCH 09/14] core/main: Don't invalidate port attributes each tick --- qtoggleserver/core/main.py | 5 ----- 1 file changed, 5 deletions(-) diff --git a/qtoggleserver/core/main.py b/qtoggleserver/core/main.py index cb1ee55e..07f8000e 100644 --- a/qtoggleserver/core/main.py +++ b/qtoggleserver/core/main.py @@ -114,7 +114,6 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> if not port.is_enabled(): continue - port.invalidate_attrs() old_value = port.get_last_read_value() if second_changed: @@ -178,10 +177,6 @@ async def handle_changes( # * ports # * time strings - # deps contain: - # * `$`-prefixed port ids - # * time strings - forced_ports = set(_force_eval_expression_ports) _force_eval_expression_ports.clear() From 1b155b15da5de489b65652d992b660b249baa365 Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Sun, 17 May 2026 23:04:39 +0300 Subject: [PATCH 10/14] core/main: Split deps into time and port categories --- qtoggleserver/core/main.py | 49 ++++++++-------------- qtoggleserver/drivers/persist/json.py | 4 +- qtoggleserver/slaves/discover/apclients.py | 2 +- tests/unit/qtoggleserver/core/test_main.py | 47 ++++++++++++--------- 4 files changed, 49 insertions(+), 53 deletions(-) diff --git a/qtoggleserver/core/main.py b/qtoggleserver/core/main.py index 07f8000e..80b29df2 100644 --- a/qtoggleserver/core/main.py +++ b/qtoggleserver/core/main.py @@ -67,8 +67,8 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> _update_lock = asyncio.Lock() async with _update_lock: - changed_set: set[core_ports.BasePort | str] = {DEP_ASAP} - value_pairs = {} + changed_time_deps: set[str] = {DEP_ASAP} + port_changed_values: dict[core_ports.BasePort, tuple[NullablePortValue, NullablePortValue]] = {} now = time.time() now_int = int(now) @@ -81,30 +81,30 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> if now_int != _last_time: _last_time = now_int second_changed = True - changed_set.add(DEP_SECOND) + changed_time_deps.add(DEP_SECOND) now_minute = now_int // 60 if now_minute != _last_minute: _last_minute = now_minute - changed_set.add(DEP_MINUTE) + changed_time_deps.add(DEP_MINUTE) now_hour = now_minute // 60 if now_hour != _last_hour: _last_hour = now_hour - changed_set.add(DEP_HOUR) + changed_time_deps.add(DEP_HOUR) now_dt = datetime.fromtimestamp(now) if now_dt.day != _last_day: _last_day = now_dt.day - changed_set.add(DEP_DAY) + changed_time_deps.add(DEP_DAY) if now_dt.month != _last_month: _last_month = now_dt.month - changed_set.add(DEP_MONTH) + changed_time_deps.add(DEP_MONTH) if now_dt.year != _last_year: _last_year = now_dt.year - changed_set.add(DEP_YEAR) + changed_time_deps.add(DEP_YEAR) all_ports = list(core_ports.get_all()) if not ports_to_read: @@ -143,10 +143,9 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> logger.debug("detected %s value change: %s -> %s", port, old_value_str, new_value_str) port.set_last_read_value(new_value) - changed_set.add(port) - value_pairs[port] = old_value, new_value + port_changed_values[port] = old_value, new_value - await handle_changes(all_ports, changed_set, value_pairs, now_ms) + await handle_changes(all_ports, changed_time_deps, port_changed_values, now_ms) sessions.update() @@ -167,39 +166,27 @@ async def update_loop() -> None: async def handle_changes( all_ports: list[core_ports.BasePort], - changed_set: set[core_ports.BasePort | str], - value_pairs: dict[core_ports.BasePort, tuple[NullablePortValue, NullablePortValue]], + changed_time_deps: set[str], + port_changed_values: dict[core_ports.BasePort, tuple[NullablePortValue, NullablePortValue]], now_ms: int, ) -> None: global _force_eval_all_expressions - # changed_set contains: - # * ports - # * time strings - forced_ports = set(_force_eval_expression_ports) _force_eval_expression_ports.clear() full_eval = _force_eval_all_expressions _force_eval_all_expressions = False - # Transform `changed_set` into a set of strings so that we can compare it with deps - changed_set_str: set[str] = set() + # Build a unified string set from time deps and port ids for dep comparison + changed_set_str: set[str] = set(changed_time_deps) - # Trigger value-change events; save persisted ports; build changed_set_str - for changed in changed_set: - if isinstance(changed, core_ports.BasePort): - port = changed - changed_set_str.add(f"${port.get_id()}") - else: # time string - changed_set_str.add(changed) - continue + # Trigger value-change events; save persisted ports; add port deps to changed_set_str + for port, (old_value, new_value) in port_changed_values.items(): + changed_set_str.add(f"${port.get_id()}") if not await port.is_internal(): - value_pair = value_pairs.get(port) - if not value_pair: - continue - await port.trigger_value_change(*value_pair) + await port.trigger_value_change(old_value, new_value) if await port.is_persisted(): port.save_asap() diff --git a/qtoggleserver/drivers/persist/json.py b/qtoggleserver/drivers/persist/json.py index 5dc70e04..e5211cee 100644 --- a/qtoggleserver/drivers/persist/json.py +++ b/qtoggleserver/drivers/persist/json.py @@ -124,7 +124,7 @@ async def insert(self, collection: str, record: Record) -> Id: else: try: self._max_ids[collection] = max(self._max_ids.get(collection, 0), int(id_)) - except (ValueError, TypeError): + except ValueError, TypeError: pass coll[id_] = record @@ -233,7 +233,7 @@ def _compute_max_ids(data: IndexedData) -> dict[str, int]: for id_ in records: try: max_id = max(max_id, int(id_)) - except (ValueError, TypeError): + except ValueError, TypeError: pass max_ids[coll] = max_id return max_ids diff --git a/qtoggleserver/slaves/discover/apclients.py b/qtoggleserver/slaves/discover/apclients.py index 1f442483..fb74d9d5 100644 --- a/qtoggleserver/slaves/discover/apclients.py +++ b/qtoggleserver/slaves/discover/apclients.py @@ -316,7 +316,7 @@ async def _query_client(ap_client: ap.APClient) -> DiscoveredDevice | None: try: attrs = await ap_client.request("GET", f"{prefix}/device") break - except (httpclient.HTTPError, json.JSONDecodeError): + except httpclient.HTTPError, json.JSONDecodeError: continue else: raise DiscoverException("Could not find device API endpoint") diff --git a/tests/unit/qtoggleserver/core/test_main.py b/tests/unit/qtoggleserver/core/test_main.py index 0bcad7db..5f236ed1 100644 --- a/tests/unit/qtoggleserver/core/test_main.py +++ b/tests/unit/qtoggleserver/core/test_main.py @@ -155,7 +155,10 @@ async def test_self_port_value_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") await handle_changes( - [mock_num_port1], changed_set={mock_num_port1}, value_pairs={mock_num_port1: (10, 20)}, now_ms=0 + [mock_num_port1], + changed_time_deps=set(), + port_changed_values={mock_num_port1: (10, 20)}, + now_ms=0, ) mock_num_port1.eval_and_push_write.assert_called_once() @@ -167,7 +170,10 @@ async def test_own_port_value_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") await handle_changes( - [mock_num_port1], changed_set={mock_num_port1}, value_pairs={mock_num_port1: (10, 20)}, now_ms=0 + [mock_num_port1], + changed_time_deps=set(), + port_changed_values={mock_num_port1: (10, 20)}, + now_ms=0, ) mock_num_port1.eval_and_push_write.assert_called_once() @@ -181,7 +187,7 @@ async def test_disabled_port_no_trigger_eval(self, mocker, mock_num_port1): (mocker.patch.object(mock_num_port1, "eval_and_push_write"),) (mocker.patch.object(mock_num_port1, "is_enabled", return_value=False),) - await handle_changes([mock_num_port1], changed_set=set(), value_pairs={}, now_ms=0) + await handle_changes([mock_num_port1], changed_time_deps=set(), port_changed_values={}, now_ms=0) mock_num_port1.eval_and_push_write.assert_not_called() async def test_asap_trigger_eval(self, mocker, mock_num_port1): @@ -191,7 +197,7 @@ async def test_asap_trigger_eval(self, mocker, mock_num_port1): mock_num_port1.set_expression("TIMEMS()") mocker.patch.object(mock_num_port1, "eval_and_push_write") - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=0) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_asap_eval_paused_no_trigger_eval(self, mocker, mock_num_port1): @@ -203,7 +209,7 @@ async def test_asap_eval_paused_no_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") e.pause_asap_eval(1000) - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=999) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=999) mock_num_port1.eval_and_push_write.assert_not_called() async def test_asap_eval_not_paused_trigger_eval(self, mocker, mock_num_port1): @@ -215,7 +221,7 @@ async def test_asap_eval_not_paused_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") e.pause_asap_eval(1000) - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=1000) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=1000) mock_num_port1.eval_and_push_write.assert_called_once() async def test_removed_port_not_evaluated(self, mocker, mock_num_port2): @@ -230,8 +236,8 @@ async def test_removed_port_not_evaluated(self, mocker, mock_num_port2): await handle_changes( list(core_ports.get_all()), - changed_set={mock_num_port2}, - value_pairs={mock_num_port2: (1, 2)}, + changed_time_deps=set(), + port_changed_values={mock_num_port2: (1, 2)}, now_ms=0, ) port.eval_and_push_write.assert_not_called() @@ -244,8 +250,8 @@ async def test_expression_set_triggers_eval_via_deps(self, mocker, mock_num_port await handle_changes( [mock_num_port1, mock_num_port2], - changed_set={mock_num_port1}, - value_pairs={mock_num_port1: (1, 2)}, + changed_time_deps=set(), + port_changed_values={mock_num_port1: (1, 2)}, now_ms=0, ) mock_num_port2.eval_and_push_write.assert_called_once() @@ -260,8 +266,8 @@ async def test_expression_cleared_stops_eval(self, mocker, mock_num_port1, mock_ await handle_changes( [mock_num_port1, mock_num_port2], - changed_set={mock_num_port2}, - value_pairs={mock_num_port2: (1, 2)}, + changed_time_deps=set(), + port_changed_values={mock_num_port2: (1, 2)}, now_ms=0, ) mock_num_port1.eval_and_push_write.assert_not_called() @@ -278,14 +284,15 @@ def reset_force_eval(self): core_main._force_eval_expression_ports.clear() async def test_forced_port_evaluated_without_matching_dep(self, mocker, mock_num_port1, mock_num_port2): - """Should evaluate a forced port even when changed_set contains no dep that port's expression uses.""" + """Should evaluate a forced port even when changed_time_deps and changed_ports contain no dep that port's + expression uses.""" mock_num_port1.set_expression("MUL($nid2, 2)") mocker.patch.object(mock_num_port1, "eval_and_push_write") force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=0) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_force_all_evaluates_all_expression_ports(self, mocker, mock_num_port1, mock_num_port2): @@ -298,7 +305,9 @@ async def test_force_all_evaluates_all_expression_ports(self, mocker, mock_num_p force_eval_expressions() - await handle_changes([mock_num_port1, mock_num_port2], changed_set=set(), value_pairs={}, now_ms=0) + await handle_changes( + [mock_num_port1, mock_num_port2], changed_time_deps=set(), port_changed_values={}, now_ms=0 + ) mock_num_port1.eval_and_push_write.assert_called_once() mock_num_port2.eval_and_push_write.assert_called_once() @@ -310,8 +319,8 @@ async def test_force_state_consumed_after_handle_changes(self, mocker, mock_num_ force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=0) - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=0) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_forced_port_bypasses_asap_pause(self, mocker, mock_num_port1): @@ -324,7 +333,7 @@ async def test_forced_port_bypasses_asap_pause(self, mocker, mock_num_port1): e.pause_asap_eval(1000) force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=999) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=999) mock_num_port1.eval_and_push_write.assert_called_once() async def test_forced_port_without_expression_not_evaluated(self, mocker, mock_num_port1): @@ -334,7 +343,7 @@ async def test_forced_port_without_expression_not_evaluated(self, mocker, mock_n force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_set={DEP_ASAP}, value_pairs={}, now_ms=0) + await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) mock_num_port1.eval_and_push_write.assert_not_called() From 63dbba609042d1c57dd293283bdae55a6ef2185b Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 18 May 2026 10:23:43 +0300 Subject: [PATCH 11/14] core/main: Reorganize module --- qtoggleserver/core/main.py | 141 ++++++------ tests/unit/qtoggleserver/core/test_main.py | 248 ++++++++++++++------- 2 files changed, 239 insertions(+), 150 deletions(-) diff --git a/qtoggleserver/core/main.py b/qtoggleserver/core/main.py index 80b29df2..b5cdc2d4 100644 --- a/qtoggleserver/core/main.py +++ b/qtoggleserver/core/main.py @@ -48,8 +48,9 @@ _update_lock: asyncio.Lock | None = None -async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> None: - from . import sessions +def _get_changed_time_deps(now_int: int) -> tuple[bool, set[str]]: + """Determine which time-unit deps have changed since the last call and update the relevant module-level tracking + variables. Returns ``(second_changed, changed_time_deps)``.""" global _last_time global _last_minute @@ -58,6 +59,44 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> global _last_week global _last_month global _last_year + + changed_time_deps: set[str] = {DEP_ASAP} + second_changed = False + + if now_int != _last_time: + _last_time = now_int + second_changed = True + changed_time_deps.add(DEP_SECOND) + + now_minute = now_int // 60 + if now_minute != _last_minute: + _last_minute = now_minute + changed_time_deps.add(DEP_MINUTE) + + now_hour = now_minute // 60 + if now_hour != _last_hour: + _last_hour = now_hour + changed_time_deps.add(DEP_HOUR) + + now_dt = datetime.fromtimestamp(now_int) + if now_dt.day != _last_day: + _last_day = now_dt.day + changed_time_deps.add(DEP_DAY) + + if now_dt.month != _last_month: + _last_month = now_dt.month + changed_time_deps.add(DEP_MONTH) + + if now_dt.year != _last_year: + _last_year = now_dt.year + changed_time_deps.add(DEP_YEAR) + + return second_changed, changed_time_deps + + +async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> None: + from . import sessions + global _update_lock if _paused: @@ -67,44 +106,19 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> _update_lock = asyncio.Lock() async with _update_lock: - changed_time_deps: set[str] = {DEP_ASAP} port_changed_values: dict[core_ports.BasePort, tuple[NullablePortValue, NullablePortValue]] = {} now = time.time() now_int = int(now) now_ms = int(now * 1000) - # Determine which time units have changed since last update. But only do this during regular port reading - # update calls, when no specific `ports_to_read` are supplied. - second_changed = False if not ports_to_read: - if now_int != _last_time: - _last_time = now_int - second_changed = True - changed_time_deps.add(DEP_SECOND) - - now_minute = now_int // 60 - if now_minute != _last_minute: - _last_minute = now_minute - changed_time_deps.add(DEP_MINUTE) - - now_hour = now_minute // 60 - if now_hour != _last_hour: - _last_hour = now_hour - changed_time_deps.add(DEP_HOUR) - - now_dt = datetime.fromtimestamp(now) - if now_dt.day != _last_day: - _last_day = now_dt.day - changed_time_deps.add(DEP_DAY) - - if now_dt.month != _last_month: - _last_month = now_dt.month - changed_time_deps.add(DEP_MONTH) - - if now_dt.year != _last_year: - _last_year = now_dt.year - changed_time_deps.add(DEP_YEAR) + second_changed, changes = _get_changed_time_deps(now_int) + else: + # When `ports_to_read` are given, this call is made from outside of the main update loop. Don't touch + # time-related deps unless called from main update loop. + second_changed = False + changes = {DEP_ASAP} all_ports = list(core_ports.get_all()) if not ports_to_read: @@ -145,31 +159,22 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> port.set_last_read_value(new_value) port_changed_values[port] = old_value, new_value - await handle_changes(all_ports, changed_time_deps, port_changed_values, now_ms) + # Trigger value-change events; save persisted ports; add port deps to changes + for port, (old_value, new_value) in port_changed_values.items(): + changes.add(f"${port.get_id()}") - sessions.update() + if not await port.is_internal(): + await port.trigger_value_change(old_value, new_value) + if await port.is_persisted(): + port.save_asap() -async def update_loop() -> None: - while True: - try: - try: - if _ready: - await read_ports() - except Exception as e: - logger.error("update failed: %s", e, exc_info=True) - await asyncio.sleep(settings.core.tick_interval / 1000.0) - except asyncio.CancelledError: - logger.debug("update task cancelled") - break + await _eval_changed_expressions(changes, now_ms) + + sessions.update() -async def handle_changes( - all_ports: list[core_ports.BasePort], - changed_time_deps: set[str], - port_changed_values: dict[core_ports.BasePort, tuple[NullablePortValue, NullablePortValue]], - now_ms: int, -) -> None: +async def _eval_changed_expressions(changed_set_str: set[str], now_ms: int) -> None: global _force_eval_all_expressions forced_ports = set(_force_eval_expression_ports) @@ -178,24 +183,11 @@ async def handle_changes( full_eval = _force_eval_all_expressions _force_eval_all_expressions = False - # Build a unified string set from time deps and port ids for dep comparison - changed_set_str: set[str] = set(changed_time_deps) - - # Trigger value-change events; save persisted ports; add port deps to changed_set_str - for port, (old_value, new_value) in port_changed_values.items(): - changed_set_str.add(f"${port.get_id()}") - - if not await port.is_internal(): - await port.trigger_value_change(old_value, new_value) - - if await port.is_persisted(): - port.save_asap() - eval_context: EvalContext | None = None # Reevaluate all port expressions depending on changed set if full_eval: - ports_to_eval = all_ports + ports_to_eval = list(core_ports.get_all()) else: deps_map = expressions_utils.get_deps_map() ports_to_eval: set[core_ports.BasePort] = set(forced_ports) @@ -217,12 +209,27 @@ async def handle_changes( if changed_deps == {DEP_ASAP} and expression.is_asap_eval_paused(now_ms): continue + # Build context lazily so we skip it entirely when all ports are filtered out above if not eval_context: eval_context = await expressions_utils.build_context(now_ms) await port.eval_and_push_write(eval_context) +async def update_loop() -> None: + while True: + try: + try: + if _ready: + await read_ports() + except Exception as e: + logger.error("update failed: %s", e, exc_info=True) + await asyncio.sleep(settings.core.tick_interval / 1000.0) + except asyncio.CancelledError: + logger.debug("update task cancelled") + break + + def force_eval_expressions(port: core_ports.BasePort | None = None) -> None: global _force_eval_all_expressions diff --git a/tests/unit/qtoggleserver/core/test_main.py b/tests/unit/qtoggleserver/core/test_main.py index 5f236ed1..2393b6a9 100644 --- a/tests/unit/qtoggleserver/core/test_main.py +++ b/tests/unit/qtoggleserver/core/test_main.py @@ -7,121 +7,216 @@ from qtoggleserver.core import main as core_main from qtoggleserver.core import ports as core_ports from qtoggleserver.core.expressions import DEP_ASAP, DEP_DAY, DEP_HOUR, DEP_MINUTE, DEP_MONTH, DEP_SECOND, DEP_YEAR -from qtoggleserver.core.main import force_eval_expressions, handle_changes, pause, read_ports, resume +from qtoggleserver.core.main import _eval_changed_expressions, force_eval_expressions, pause, read_ports, resume from tests.unit.qtoggleserver.mock.ports import MockNumberPort +def _ts(dt: datetime) -> int: + return int(dt.timestamp()) + + +class TestGetChangedTimeDeps: + @pytest.fixture(autouse=True) + def reset_time_globals(self): + """Reset all _last_* module globals to 0 before and after each test.""" + for attr in ( + "_last_time", + "_last_minute", + "_last_hour", + "_last_day", + "_last_week", + "_last_month", + "_last_year", + ): + setattr(core_main, attr, 0) + yield + for attr in ( + "_last_time", + "_last_minute", + "_last_hour", + "_last_day", + "_last_week", + "_last_month", + "_last_year", + ): + setattr(core_main, attr, 0) + + def test_same_second_returns_asap_only(self): + """Calling with the same now_int twice should return second_changed=False and only DEP_ASAP.""" + t = _ts(datetime(2024, 3, 15, 10, 30, 30)) + core_main._get_changed_time_deps(t) # prime + second_changed, changed = core_main._get_changed_time_deps(t) + assert second_changed is False + assert changed == {DEP_ASAP} + + def test_second_changes(self): + """A new second within the same minute should add DEP_SECOND.""" + t0 = _ts(datetime(2024, 3, 15, 10, 30, 30)) + t1 = _ts(datetime(2024, 3, 15, 10, 30, 31)) + core_main._get_changed_time_deps(t0) + second_changed, changed = core_main._get_changed_time_deps(t1) + assert second_changed is True + assert changed == {DEP_ASAP, DEP_SECOND} + + def test_minute_changes(self): + """A new minute within the same hour should add DEP_SECOND and DEP_MINUTE.""" + t0 = _ts(datetime(2024, 3, 15, 10, 30, 30)) + t1 = _ts(datetime(2024, 3, 15, 10, 31, 0)) + core_main._get_changed_time_deps(t0) + second_changed, changed = core_main._get_changed_time_deps(t1) + assert second_changed is True + assert changed == {DEP_ASAP, DEP_SECOND, DEP_MINUTE} + + def test_hour_changes(self): + """A new hour within the same day should add DEP_SECOND, DEP_MINUTE and DEP_HOUR.""" + t0 = _ts(datetime(2024, 3, 15, 10, 30, 30)) + t1 = _ts(datetime(2024, 3, 15, 11, 0, 0)) + core_main._get_changed_time_deps(t0) + second_changed, changed = core_main._get_changed_time_deps(t1) + assert second_changed is True + assert changed == {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR} + + def test_day_changes(self): + """A new day within the same month should add DEP_SECOND through DEP_DAY.""" + t0 = _ts(datetime(2024, 3, 15, 10, 30, 30)) + t1 = _ts(datetime(2024, 3, 16, 10, 30, 0)) + core_main._get_changed_time_deps(t0) + second_changed, changed = core_main._get_changed_time_deps(t1) + assert second_changed is True + assert changed == {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY} + + def test_month_changes(self): + """A new month within the same year should add DEP_SECOND through DEP_MONTH.""" + t0 = _ts(datetime(2024, 1, 31, 10, 30, 30)) + t1 = _ts(datetime(2024, 2, 1, 10, 30, 0)) + core_main._get_changed_time_deps(t0) + second_changed, changed = core_main._get_changed_time_deps(t1) + assert second_changed is True + assert changed == {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH} + + def test_year_changes(self): + """A new year should add all time deps including DEP_YEAR.""" + t0 = _ts(datetime(2023, 12, 31, 10, 30, 30)) + t1 = _ts(datetime(2024, 1, 1, 10, 30, 0)) + core_main._get_changed_time_deps(t0) + second_changed, changed = core_main._get_changed_time_deps(t1) + assert second_changed is True + assert changed == {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH, DEP_YEAR} + + def test_globals_updated(self): + """After a call, all relevant _last_* globals should reflect the new time.""" + t = _ts(datetime(2024, 3, 15, 10, 30, 30)) + dt = datetime.fromtimestamp(t) + core_main._get_changed_time_deps(t) + assert core_main._last_time == t + assert core_main._last_minute == t // 60 + assert core_main._last_hour == t // 3600 + assert core_main._last_day == dt.day + assert core_main._last_month == dt.month + assert core_main._last_year == dt.year + + class TestReadPorts: async def test_change_time_asap(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should call `handle_changes` with {DEP_ASAP}, regardless of time changes.""" + """Should call `_eval_changed_expressions` with {DEP_ASAP}, regardless of time changes.""" freezer.move_to(dummy_utc_datetime) await read_ports() - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() - spy_handle_value_changes.assert_called_once_with([mock_num_port1], {DEP_ASAP}, {}, int(time.time() * 1000)) + spy_handle_value_changes.assert_called_once_with({DEP_ASAP}, int(time.time() * 1000)) async def test_change_time_second(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should call `handle_changes` with {DEP_ASAP, DEP_SECOND} when second changes.""" + """Should call `_eval_changed_expressions` with {DEP_ASAP, DEP_SECOND} when second changes.""" freezer.move_to(dummy_utc_datetime) await read_ports() freezer.move_to(dummy_utc_datetime + timedelta(seconds=1)) - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() - spy_handle_value_changes.assert_called_once_with( - [mock_num_port1], {DEP_ASAP, DEP_SECOND}, {}, int(time.time() * 1000) - ) + spy_handle_value_changes.assert_called_once_with({DEP_ASAP, DEP_SECOND}, int(time.time() * 1000)) async def test_change_time_minute(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should call `handle_changes` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE} when minute changes.""" + """Should call `_eval_changed_expressions` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE} when minute changes.""" freezer.move_to(dummy_utc_datetime) await read_ports() freezer.move_to(dummy_utc_datetime + timedelta(minutes=1)) - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() - spy_handle_value_changes.assert_called_once_with( - [mock_num_port1], {DEP_ASAP, DEP_SECOND, DEP_MINUTE}, {}, int(time.time() * 1000) - ) + spy_handle_value_changes.assert_called_once_with({DEP_ASAP, DEP_SECOND, DEP_MINUTE}, int(time.time() * 1000)) async def test_change_time_hour(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should call `handle_changes` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR} when hour changes.""" + """Should call `_eval_changed_expressions` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR} whenever hour + changes.""" freezer.move_to(dummy_utc_datetime) await read_ports() freezer.move_to(dummy_utc_datetime + timedelta(hours=1)) - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy_handle_value_changes.assert_called_once_with( - [mock_num_port1], {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR}, {}, int(time.time() * 1000) + {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR}, int(time.time() * 1000) ) async def test_change_time_day(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should call `handle_changes` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY} when day + """Should call `_eval_changed_expressions` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY} when day changes.""" freezer.move_to(datetime(2019, 1, 30, 23, 30, 30)) await read_ports() freezer.move_to(datetime(2019, 1, 31, 0, 0, 0)) - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy_handle_value_changes.assert_called_once_with( - [mock_num_port1], {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY}, {}, int(time.time() * 1000) + {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY}, int(time.time() * 1000) ) async def test_change_time_month(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should call `handle_changes` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH} when - month changes.""" + """Should call `_eval_changed_expressions` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH} + when month changes.""" freezer.move_to(datetime(2019, 1, 31, 23, 30, 30)) await read_ports() freezer.move_to(datetime(2019, 2, 1, 0, 0, 0)) - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy_handle_value_changes.assert_called_once_with( - [mock_num_port1], {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH}, - {}, int(time.time() * 1000), ) async def test_change_time_year(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should call `handle_changes` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH, + """Should call `_eval_changed_expressions` with {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH, DEP_YEAR} when year changes.""" freezer.move_to(datetime(2019, 12, 31, 23, 30, 30)) await read_ports() freezer.move_to(datetime(2020, 1, 1, 0, 0, 0)) - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy_handle_value_changes.assert_called_once_with( - [mock_num_port1], {DEP_ASAP, DEP_SECOND, DEP_MINUTE, DEP_HOUR, DEP_DAY, DEP_MONTH, DEP_YEAR}, - {}, int(time.time() * 1000), ) async def test_ports_to_read_specific_ports( self, freezer, mocker, mock_num_port1, mock_num_port2, dummy_utc_datetime ): - """Should only read specified ports when `ports_to_read` is provided, while passing all ports to - handle_changes.""" + """Should only read specified ports when `ports_to_read` is provided.""" freezer.move_to(dummy_utc_datetime) await read_ports() - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports(ports_to_read=[mock_num_port2]) - spy_handle_value_changes.assert_called_once_with( - [mock_num_port1, mock_num_port2], {DEP_ASAP}, {}, int(time.time() * 1000) - ) + spy_handle_value_changes.assert_called_once_with({DEP_ASAP}, int(time.time() * 1000)) async def test_ports_to_read_skip_time_changes(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): """Should not add time dependencies when `ports_to_read` is provided, even if time changes.""" @@ -130,20 +225,20 @@ async def test_ports_to_read_skip_time_changes(self, freezer, mocker, mock_num_p await read_ports() freezer.move_to(dummy_utc_datetime + timedelta(seconds=1)) - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports(ports_to_read=[mock_num_port1]) - spy_handle_value_changes.assert_called_once_with([mock_num_port1], {DEP_ASAP}, {}, int(time.time() * 1000)) + spy_handle_value_changes.assert_called_once_with({DEP_ASAP}, int(time.time() * 1000)) async def test_ports_to_read_empty_list(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): """Should handle empty ports list gracefully when `ports_to_read` is an empty list.""" freezer.move_to(dummy_utc_datetime) await read_ports() - spy_handle_value_changes = mocker.patch("qtoggleserver.core.main.handle_changes") + spy_handle_value_changes = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports(ports_to_read=[]) - spy_handle_value_changes.assert_called_once_with([mock_num_port1], {DEP_ASAP}, {}, int(time.time() * 1000)) + spy_handle_value_changes.assert_called_once_with({DEP_ASAP}, int(time.time() * 1000)) class TestHandleChanges: @@ -154,10 +249,8 @@ async def test_self_port_value_trigger_eval(self, mocker, mock_num_port1): mock_num_port1.set_expression("MUL($, 2)") mocker.patch.object(mock_num_port1, "eval_and_push_write") - await handle_changes( - [mock_num_port1], - changed_time_deps=set(), - port_changed_values={mock_num_port1: (10, 20)}, + await _eval_changed_expressions( + changed_set_str={"$nid1"}, now_ms=0, ) mock_num_port1.eval_and_push_write.assert_called_once() @@ -169,10 +262,8 @@ async def test_own_port_value_trigger_eval(self, mocker, mock_num_port1): mock_num_port1.set_expression("MUL($nid1, 2)") mocker.patch.object(mock_num_port1, "eval_and_push_write") - await handle_changes( - [mock_num_port1], - changed_time_deps=set(), - port_changed_values={mock_num_port1: (10, 20)}, + await _eval_changed_expressions( + changed_set_str={"$nid1"}, now_ms=0, ) mock_num_port1.eval_and_push_write.assert_called_once() @@ -187,7 +278,7 @@ async def test_disabled_port_no_trigger_eval(self, mocker, mock_num_port1): (mocker.patch.object(mock_num_port1, "eval_and_push_write"),) (mocker.patch.object(mock_num_port1, "is_enabled", return_value=False),) - await handle_changes([mock_num_port1], changed_time_deps=set(), port_changed_values={}, now_ms=0) + await _eval_changed_expressions(changed_set_str=set(), now_ms=0) mock_num_port1.eval_and_push_write.assert_not_called() async def test_asap_trigger_eval(self, mocker, mock_num_port1): @@ -197,7 +288,7 @@ async def test_asap_trigger_eval(self, mocker, mock_num_port1): mock_num_port1.set_expression("TIMEMS()") mocker.patch.object(mock_num_port1, "eval_and_push_write") - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_asap_eval_paused_no_trigger_eval(self, mocker, mock_num_port1): @@ -209,7 +300,7 @@ async def test_asap_eval_paused_no_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") e.pause_asap_eval(1000) - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=999) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=999) mock_num_port1.eval_and_push_write.assert_not_called() async def test_asap_eval_not_paused_trigger_eval(self, mocker, mock_num_port1): @@ -221,7 +312,7 @@ async def test_asap_eval_not_paused_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") e.pause_asap_eval(1000) - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=1000) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=1000) mock_num_port1.eval_and_push_write.assert_called_once() async def test_removed_port_not_evaluated(self, mocker, mock_num_port2): @@ -234,10 +325,8 @@ async def test_removed_port_not_evaluated(self, mocker, mock_num_port2): await port.remove(persisted_data=False) - await handle_changes( - list(core_ports.get_all()), - changed_time_deps=set(), - port_changed_values={mock_num_port2: (1, 2)}, + await _eval_changed_expressions( + changed_set_str={"$nid2"}, now_ms=0, ) port.eval_and_push_write.assert_not_called() @@ -248,10 +337,8 @@ async def test_expression_set_triggers_eval_via_deps(self, mocker, mock_num_port mock_num_port2.set_expression("MUL($nid1, 2)") mocker.patch.object(mock_num_port2, "eval_and_push_write") - await handle_changes( - [mock_num_port1, mock_num_port2], - changed_time_deps=set(), - port_changed_values={mock_num_port1: (1, 2)}, + await _eval_changed_expressions( + changed_set_str={"$nid1"}, now_ms=0, ) mock_num_port2.eval_and_push_write.assert_called_once() @@ -264,10 +351,8 @@ async def test_expression_cleared_stops_eval(self, mocker, mock_num_port1, mock_ await mock_num_port1.attr_set_expression("") mocker.patch.object(mock_num_port1, "eval_and_push_write") - await handle_changes( - [mock_num_port1, mock_num_port2], - changed_time_deps=set(), - port_changed_values={mock_num_port2: (1, 2)}, + await _eval_changed_expressions( + changed_set_str={"$nid2"}, now_ms=0, ) mock_num_port1.eval_and_push_write.assert_not_called() @@ -284,15 +369,14 @@ def reset_force_eval(self): core_main._force_eval_expression_ports.clear() async def test_forced_port_evaluated_without_matching_dep(self, mocker, mock_num_port1, mock_num_port2): - """Should evaluate a forced port even when changed_time_deps and changed_ports contain no dep that port's - expression uses.""" + """Should evaluate a forced port even when changed_set_str contains no dep that port's expression uses.""" mock_num_port1.set_expression("MUL($nid2, 2)") mocker.patch.object(mock_num_port1, "eval_and_push_write") force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_force_all_evaluates_all_expression_ports(self, mocker, mock_num_port1, mock_num_port2): @@ -305,22 +389,20 @@ async def test_force_all_evaluates_all_expression_ports(self, mocker, mock_num_p force_eval_expressions() - await handle_changes( - [mock_num_port1, mock_num_port2], changed_time_deps=set(), port_changed_values={}, now_ms=0 - ) + await _eval_changed_expressions(changed_set_str=set(), now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() mock_num_port2.eval_and_push_write.assert_called_once() - async def test_force_state_consumed_after_handle_changes(self, mocker, mock_num_port1, mock_num_port2): - """Should not re-evaluate a forced port on a subsequent handle_changes call without re-forcing.""" + async def test_force_state_consumed_after_eval_changed_expressions(self, mocker, mock_num_port1, mock_num_port2): + """Should not re-evaluate a forced port on a subsequent _eval_changed_expressions call without re-forcing.""" mock_num_port1.set_expression("MUL($nid2, 2)") mocker.patch.object(mock_num_port1, "eval_and_push_write") force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_forced_port_bypasses_asap_pause(self, mocker, mock_num_port1): @@ -333,7 +415,7 @@ async def test_forced_port_bypasses_asap_pause(self, mocker, mock_num_port1): e.pause_asap_eval(1000) force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=999) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=999) mock_num_port1.eval_and_push_write.assert_called_once() async def test_forced_port_without_expression_not_evaluated(self, mocker, mock_num_port1): @@ -343,7 +425,7 @@ async def test_forced_port_without_expression_not_evaluated(self, mocker, mock_n force_eval_expressions(mock_num_port1) - await handle_changes([mock_num_port1], changed_time_deps={DEP_ASAP}, port_changed_values={}, now_ms=0) + await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_not_called() @@ -356,19 +438,19 @@ def reset_paused(self): core_main._paused = False async def test_resume_allows_read_ports(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """Should allow read_ports() to call handle_changes after resume().""" + """Should allow read_ports() to call _eval_changed_expressions after resume().""" freezer.move_to(dummy_utc_datetime) await read_ports() # prime last-time state - spy = mocker.patch("qtoggleserver.core.main.handle_changes") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") resume() await read_ports() spy.assert_called_once() async def test_pause_skips_read_ports(self, mocker, mock_num_port1): - """Should skip handle_changes call in read_ports() after pause().""" + """Should skip _eval_changed_expressions call in read_ports() after pause().""" - spy = mocker.patch("qtoggleserver.core.main.handle_changes") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") pause() await read_ports() spy.assert_not_called() @@ -380,7 +462,7 @@ async def test_resume_from_paused_state(self, freezer, mocker, mock_num_port1, d await read_ports() # prime last-time state pause() resume() - spy = mocker.patch("qtoggleserver.core.main.handle_changes") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy.assert_called_once() @@ -391,14 +473,14 @@ async def test_resume_idempotent(self, freezer, mocker, mock_num_port1, dummy_ut await read_ports() # prime last-time state resume() resume() - spy = mocker.patch("qtoggleserver.core.main.handle_changes") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy.assert_called_once() async def test_pause_idempotent(self, mocker, mock_num_port1): """Calling pause() multiple times should keep read_ports() skipped.""" - spy = mocker.patch("qtoggleserver.core.main.handle_changes") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") pause() pause() await read_ports() @@ -413,6 +495,6 @@ async def test_pause_resume_cycle(self, freezer, mocker, mock_num_port1, dummy_u resume() pause() resume() - spy = mocker.patch("qtoggleserver.core.main.handle_changes") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy.assert_called_once() From ae3c3249872b743a6dbb9eed46612fbfe29b9ebf Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 18 May 2026 22:37:23 +0300 Subject: [PATCH 12/14] core/main: Listen for events and eval expressions depending on attrs --- qtoggleserver/core/main.py | 34 +++- tests/unit/qtoggleserver/core/test_main.py | 199 +++++++++++++++++++-- 2 files changed, 215 insertions(+), 18 deletions(-) diff --git a/qtoggleserver/core/main.py b/qtoggleserver/core/main.py index b5cdc2d4..1bae2fec 100644 --- a/qtoggleserver/core/main.py +++ b/qtoggleserver/core/main.py @@ -5,6 +5,7 @@ from datetime import datetime from qtoggleserver.conf import settings +from qtoggleserver.core import events as core_events from qtoggleserver.core import ports as core_ports from qtoggleserver.core.expressions import ( DEP_ASAP, @@ -45,9 +46,30 @@ _force_eval_expression_ports: set[core_ports.BasePort] = set() _force_eval_all_expressions: bool = False _ports_with_read_error = timedset.TimedSet(_PORT_READ_ERROR_RETRY_INTERVAL) +_pending_attr_changes: set[str] = set() _update_lock: asyncio.Lock | None = None +class _AttrChangeHandler(core_events.Handler): + """Flags port/device attribute dep changes for re-evaluation on the next read_ports tick.""" + + FIRE_AND_FORGET = False + + async def handle_event(self, event: core_events.Event) -> None: + if isinstance(event, (core_events.PortAdd, core_events.PortRemove, core_events.PortUpdate)): + _pending_attr_changes.add(f"${event.get_port().get_id()}:") + elif isinstance(event, core_events.DeviceUpdate): + _pending_attr_changes.add("#:") + else: + from qtoggleserver.slaves import events as slaves_events + + if isinstance( + event, + (slaves_events.SlaveDeviceAdd, slaves_events.SlaveDeviceRemove, slaves_events.SlaveDeviceUpdate), + ): + _pending_attr_changes.add(f"#{event.get_slave().get_name()}:") + + def _get_changed_time_deps(now_int: int) -> tuple[bool, set[str]]: """Determine which time-unit deps have changed since the last call and update the relevant module-level tracking variables. Returns ``(second_changed, changed_time_deps)``.""" @@ -120,6 +142,11 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> second_changed = False changes = {DEP_ASAP} + # Drain any attribute dep changes (port or device) flagged by the event handler since the last tick + if _pending_attr_changes: + changes.update(_pending_attr_changes) + _pending_attr_changes.clear() + all_ports = list(core_ports.get_all()) if not ports_to_read: ports_to_read = all_ports @@ -174,7 +201,7 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> sessions.update() -async def _eval_changed_expressions(changed_set_str: set[str], now_ms: int) -> None: +async def _eval_changed_expressions(changes: set[str], now_ms: int) -> None: global _force_eval_all_expressions forced_ports = set(_force_eval_expression_ports) @@ -191,7 +218,7 @@ async def _eval_changed_expressions(changed_set_str: set[str], now_ms: int) -> N else: deps_map = expressions_utils.get_deps_map() ports_to_eval: set[core_ports.BasePort] = set(forced_ports) - for dep in changed_set_str: + for dep in changes: for port in deps_map.get(dep, []): ports_to_eval.add(port) @@ -205,7 +232,7 @@ async def _eval_changed_expressions(changed_set_str: set[str], now_ms: int) -> N if not full_eval and port not in forced_ports: deps: set[str] = expression.get_deps() - changed_deps = deps & changed_set_str + changed_deps = deps & changes if changed_deps == {DEP_ASAP} and expression.is_asap_eval_paused(now_ms): continue @@ -278,6 +305,7 @@ async def init() -> None: loop = asyncio.get_running_loop() + core_events.register_handler(_AttrChangeHandler()) force_eval_expressions() _update_loop_task = loop.create_task(update_loop()) diff --git a/tests/unit/qtoggleserver/core/test_main.py b/tests/unit/qtoggleserver/core/test_main.py index 2393b6a9..105f43d5 100644 --- a/tests/unit/qtoggleserver/core/test_main.py +++ b/tests/unit/qtoggleserver/core/test_main.py @@ -1,16 +1,25 @@ import time from datetime import datetime, timedelta +from unittest.mock import MagicMock import pytest +from qtoggleserver.core import events as core_events from qtoggleserver.core import main as core_main from qtoggleserver.core import ports as core_ports from qtoggleserver.core.expressions import DEP_ASAP, DEP_DAY, DEP_HOUR, DEP_MINUTE, DEP_MONTH, DEP_SECOND, DEP_YEAR from qtoggleserver.core.main import _eval_changed_expressions, force_eval_expressions, pause, read_ports, resume +from qtoggleserver.slaves import events as slaves_events from tests.unit.qtoggleserver.mock.ports import MockNumberPort +def _make_mock_slave(name: str) -> MagicMock: + slave = MagicMock() + slave.get_name.return_value = name + return slave + + def _ts(dt: datetime) -> int: return int(dt.timestamp()) @@ -250,7 +259,7 @@ async def test_self_port_value_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") await _eval_changed_expressions( - changed_set_str={"$nid1"}, + changes={"$nid1"}, now_ms=0, ) mock_num_port1.eval_and_push_write.assert_called_once() @@ -263,7 +272,7 @@ async def test_own_port_value_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") await _eval_changed_expressions( - changed_set_str={"$nid1"}, + changes={"$nid1"}, now_ms=0, ) mock_num_port1.eval_and_push_write.assert_called_once() @@ -278,7 +287,7 @@ async def test_disabled_port_no_trigger_eval(self, mocker, mock_num_port1): (mocker.patch.object(mock_num_port1, "eval_and_push_write"),) (mocker.patch.object(mock_num_port1, "is_enabled", return_value=False),) - await _eval_changed_expressions(changed_set_str=set(), now_ms=0) + await _eval_changed_expressions(changes=set(), now_ms=0) mock_num_port1.eval_and_push_write.assert_not_called() async def test_asap_trigger_eval(self, mocker, mock_num_port1): @@ -288,7 +297,7 @@ async def test_asap_trigger_eval(self, mocker, mock_num_port1): mock_num_port1.set_expression("TIMEMS()") mocker.patch.object(mock_num_port1, "eval_and_push_write") - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_asap_eval_paused_no_trigger_eval(self, mocker, mock_num_port1): @@ -300,7 +309,7 @@ async def test_asap_eval_paused_no_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") e.pause_asap_eval(1000) - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=999) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=999) mock_num_port1.eval_and_push_write.assert_not_called() async def test_asap_eval_not_paused_trigger_eval(self, mocker, mock_num_port1): @@ -312,7 +321,7 @@ async def test_asap_eval_not_paused_trigger_eval(self, mocker, mock_num_port1): mocker.patch.object(mock_num_port1, "eval_and_push_write") e.pause_asap_eval(1000) - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=1000) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=1000) mock_num_port1.eval_and_push_write.assert_called_once() async def test_removed_port_not_evaluated(self, mocker, mock_num_port2): @@ -326,7 +335,7 @@ async def test_removed_port_not_evaluated(self, mocker, mock_num_port2): await port.remove(persisted_data=False) await _eval_changed_expressions( - changed_set_str={"$nid2"}, + changes={"$nid2"}, now_ms=0, ) port.eval_and_push_write.assert_not_called() @@ -338,7 +347,7 @@ async def test_expression_set_triggers_eval_via_deps(self, mocker, mock_num_port mocker.patch.object(mock_num_port2, "eval_and_push_write") await _eval_changed_expressions( - changed_set_str={"$nid1"}, + changes={"$nid1"}, now_ms=0, ) mock_num_port2.eval_and_push_write.assert_called_once() @@ -352,7 +361,7 @@ async def test_expression_cleared_stops_eval(self, mocker, mock_num_port1, mock_ mocker.patch.object(mock_num_port1, "eval_and_push_write") await _eval_changed_expressions( - changed_set_str={"$nid2"}, + changes={"$nid2"}, now_ms=0, ) mock_num_port1.eval_and_push_write.assert_not_called() @@ -376,7 +385,7 @@ async def test_forced_port_evaluated_without_matching_dep(self, mocker, mock_num force_eval_expressions(mock_num_port1) - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_force_all_evaluates_all_expression_ports(self, mocker, mock_num_port1, mock_num_port2): @@ -389,7 +398,7 @@ async def test_force_all_evaluates_all_expression_ports(self, mocker, mock_num_p force_eval_expressions() - await _eval_changed_expressions(changed_set_str=set(), now_ms=0) + await _eval_changed_expressions(changes=set(), now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() mock_num_port2.eval_and_push_write.assert_called_once() @@ -401,8 +410,8 @@ async def test_force_state_consumed_after_eval_changed_expressions(self, mocker, force_eval_expressions(mock_num_port1) - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=0) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_called_once() async def test_forced_port_bypasses_asap_pause(self, mocker, mock_num_port1): @@ -415,7 +424,7 @@ async def test_forced_port_bypasses_asap_pause(self, mocker, mock_num_port1): e.pause_asap_eval(1000) force_eval_expressions(mock_num_port1) - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=999) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=999) mock_num_port1.eval_and_push_write.assert_called_once() async def test_forced_port_without_expression_not_evaluated(self, mocker, mock_num_port1): @@ -425,7 +434,7 @@ async def test_forced_port_without_expression_not_evaluated(self, mocker, mock_n force_eval_expressions(mock_num_port1) - await _eval_changed_expressions(changed_set_str={DEP_ASAP}, now_ms=0) + await _eval_changed_expressions(changes={DEP_ASAP}, now_ms=0) mock_num_port1.eval_and_push_write.assert_not_called() @@ -498,3 +507,163 @@ async def test_pause_resume_cycle(self, freezer, mocker, mock_num_port1, dummy_u spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() spy.assert_called_once() + + +class TestAttrChangeHandler: + @pytest.fixture(autouse=True) + def reset_pending(self): + """Clear _pending_attr_changes before and after each test.""" + core_main._pending_attr_changes.clear() + yield + core_main._pending_attr_changes.clear() + + async def test_port_update_adds_attr_dep(self, mock_num_port1): + """PortUpdate event should add `$port_id:` to _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(core_events.PortUpdate(mock_num_port1)) + assert core_main._pending_attr_changes == {"$nid1:"} + + async def test_port_add_adds_attr_dep(self, mock_num_port1): + """PortAdd event should add `$port_id:` to _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(core_events.PortAdd(mock_num_port1)) + assert core_main._pending_attr_changes == {"$nid1:"} + + async def test_port_remove_adds_attr_dep(self, mock_num_port1): + """PortRemove event should add `$port_id:` to _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(core_events.PortRemove(mock_num_port1)) + assert core_main._pending_attr_changes == {"$nid1:"} + + async def test_device_update_adds_device_dep(self): + """DeviceUpdate event should add `#:` to _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(core_events.DeviceUpdate()) + assert core_main._pending_attr_changes == {"#:"} + + async def test_value_change_ignored(self, mock_num_port1): + """ValueChange event should not modify _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(core_events.ValueChange(None, 42, mock_num_port1)) + assert core_main._pending_attr_changes == set() + + async def test_multiple_events_accumulate(self, mock_num_port1, mock_num_port2): + """Multiple port events should accumulate their dep strings.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(core_events.PortUpdate(mock_num_port1)) + await handler.handle_event(core_events.PortUpdate(mock_num_port2)) + assert core_main._pending_attr_changes == {"$nid1:", "$nid2:"} + + async def test_port_and_device_events_accumulate(self, mock_num_port1): + """Port and device events should both accumulate into _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(core_events.PortUpdate(mock_num_port1)) + await handler.handle_event(core_events.DeviceUpdate()) + assert core_main._pending_attr_changes == {"$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.""" + + freezer.move_to(dummy_utc_datetime) + await read_ports() # prime time state + + core_main._pending_attr_changes.add("$nid1:") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") + await read_ports() + + call_changes = spy.call_args[0][0] + assert "$nid1:" in call_changes + + async def test_read_ports_drains_device_dep(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): + """read_ports() should include `#:` in changes when a DeviceUpdate was pending.""" + + freezer.move_to(dummy_utc_datetime) + await read_ports() # prime time state + + core_main._pending_attr_changes.add("#:") + spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") + await read_ports() + + call_changes = spy.call_args[0][0] + assert "#:" in call_changes + + async def test_read_ports_clears_pending_after_drain(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): + """_pending_attr_changes must be empty after read_ports() drains it.""" + + freezer.move_to(dummy_utc_datetime) + await read_ports() # prime time state + + core_main._pending_attr_changes.add("$nid1:") + mocker.patch("qtoggleserver.core.main._eval_changed_expressions") + await read_ports() + + assert core_main._pending_attr_changes == set() + + async def test_port_attr_dep_triggers_expression_eval(self, mocker, mock_num_port1, mock_num_port2): + """A port whose expression references $nid1: should be evaluated when `$nid1:` is in changes.""" + + mock_num_port2.set_writable(True) + mock_num_port2.set_expression("$nid1:enabled") + mocker.patch.object(mock_num_port2, "eval_and_push_write") + + await _eval_changed_expressions(changes={"$nid1:"}, 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.""" + + mock_num_port1.set_writable(True) + mock_num_port1.set_expression("#:name") + mocker.patch.object(mock_num_port1, "eval_and_push_write") + + await _eval_changed_expressions(changes={"#:"}, now_ms=0) + + mock_num_port1.eval_and_push_write.assert_called_once() + + async def test_slave_device_update_adds_slave_dep(self): + """SlaveDeviceUpdate event should add `#slave_name:` to _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1"))) + assert core_main._pending_attr_changes == {"#slave1:"} + + async def test_slave_device_add_adds_slave_dep(self): + """SlaveDeviceAdd event should add `#slave_name:` to _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(slaves_events.SlaveDeviceAdd(_make_mock_slave("slave1"))) + assert core_main._pending_attr_changes == {"#slave1:"} + + async def test_slave_device_remove_adds_slave_dep(self): + """SlaveDeviceRemove event should add `#slave_name:` to _pending_attr_changes.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(slaves_events.SlaveDeviceRemove(_make_mock_slave("slave1"))) + assert core_main._pending_attr_changes == {"#slave1:"} + + async def test_multiple_slave_events_accumulate(self): + """Multiple slave device events should accumulate distinct dep strings.""" + + handler = core_main._AttrChangeHandler() + await handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1"))) + await handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave2"))) + assert core_main._pending_attr_changes == {"#slave1:", "#slave2:"} + + async def test_slave_dep_triggers_expression_eval(self, mocker, mock_num_port1): + """A port whose expression references #slave1:attr should be evaluated when `#slave1:` is in changes.""" + + mock_num_port1.set_writable(True) + mock_num_port1.set_expression("#slave1:enabled") + mocker.patch.object(mock_num_port1, "eval_and_push_write") + + await _eval_changed_expressions(changes={"#slave1:"}, now_ms=0) + + mock_num_port1.eval_and_push_write.assert_called_once() From 4883ebbe236aa2aa9c4780b43960372ac0041033 Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 18 May 2026 22:51:59 +0300 Subject: [PATCH 13/14] core/main: Fix default argument kinds list --- .../core/expressions/functions/__init__.py | 5 ++-- .../core/expressions/test_parse.py | 26 +++++++++++++++++++ 2 files changed, 29 insertions(+), 2 deletions(-) diff --git a/qtoggleserver/core/expressions/functions/__init__.py b/qtoggleserver/core/expressions/functions/__init__.py index 0e647beb..ffde16f1 100644 --- a/qtoggleserver/core/expressions/functions/__init__.py +++ b/qtoggleserver/core/expressions/functions/__init__.py @@ -6,8 +6,9 @@ from .. import DEP_ASAP, exceptions, parse from ..base import EvalContext, EvalResult, Expression, Role +from ..devices import DeviceAttr from ..literalvalues import LiteralValue -from ..ports import PortValue +from ..ports import PortAttr, PortValue FUNCTIONS = {} @@ -73,7 +74,7 @@ def validate_arg_kinds(cls, args: list[Expression], pos_list: list[int]) -> None try: kind = cls.ARG_KINDS[i] except IndexError: - kind = (LiteralValue, PortValue, Function) + kind = (LiteralValue, PortValue, Function, PortAttr, DeviceAttr) if not isinstance(arg, kind): raise exceptions.InvalidArgumentKind(cls.NAME, pos_list[i], i + 1) diff --git a/tests/unit/qtoggleserver/core/expressions/test_parse.py b/tests/unit/qtoggleserver/core/expressions/test_parse.py index 37f1c5d0..385df78a 100644 --- a/tests/unit/qtoggleserver/core/expressions/test_parse.py +++ b/tests/unit/qtoggleserver/core/expressions/test_parse.py @@ -153,3 +153,29 @@ async def test_parse_missing_attr_prefix(): async def test_parse_empty(): with pytest.raises(EmptyExpression): parse(None, " ", role=Role.VALUE) + + +async def test_parse_port_attr_as_function_arg(mock_num_port1): + """A port-attribute expression ($port_id:attr) must be accepted as a function argument.""" + context = EvalContext( + port_values={"nid1": 0}, + port_attrs={"nid1": {"enabled": True}}, + now_ms=0, + ) + + e = parse("nid1", "ADD($nid1:enabled, 2)", role=Role.VALUE) + # enabled=True → 1, + 2 = 3 + assert await e.eval(context) == 3 + + +async def test_parse_device_attr_as_function_arg(): + """A device-attribute expression (#:attr) must be accepted as a function argument.""" + context = EvalContext( + port_values={}, + device_attrs={"name": "my_device"}, # non-numeric → int(bool(...)) = 1 + now_ms=0, + ) + + e = parse(None, "ADD(#:name, 2)", role=Role.VALUE) + # "my_device" is truthy → 1, + 2 = 3 + assert await e.eval(context) == 3 From 380bfbd0752290ed1ab6d2482343dd75c0e45f17 Mon Sep 17 00:00:00 2001 From: Calin Crisan Date: Mon, 18 May 2026 23:12:14 +0300 Subject: [PATCH 14/14] core/main: Move pending attr change logic to main.utils --- qtoggleserver/core/main.py | 35 +++------- qtoggleserver/utils/main.py | 39 +++++++++++ tests/unit/qtoggleserver/core/test_main.py | 77 ++++++++++------------ 3 files changed, 80 insertions(+), 71 deletions(-) create mode 100644 qtoggleserver/utils/main.py diff --git a/qtoggleserver/core/main.py b/qtoggleserver/core/main.py index 1bae2fec..dcc67c47 100644 --- a/qtoggleserver/core/main.py +++ b/qtoggleserver/core/main.py @@ -21,6 +21,7 @@ from qtoggleserver.utils import expressions as expressions_utils from qtoggleserver.utils import json as json_utils from qtoggleserver.utils import logging as logging_utils +from qtoggleserver.utils import main as main_utils from qtoggleserver.utils import timedset @@ -46,28 +47,8 @@ _force_eval_expression_ports: set[core_ports.BasePort] = set() _force_eval_all_expressions: bool = False _ports_with_read_error = timedset.TimedSet(_PORT_READ_ERROR_RETRY_INTERVAL) -_pending_attr_changes: set[str] = set() _update_lock: asyncio.Lock | None = None - - -class _AttrChangeHandler(core_events.Handler): - """Flags port/device attribute dep changes for re-evaluation on the next read_ports tick.""" - - FIRE_AND_FORGET = False - - async def handle_event(self, event: core_events.Event) -> None: - if isinstance(event, (core_events.PortAdd, core_events.PortRemove, core_events.PortUpdate)): - _pending_attr_changes.add(f"${event.get_port().get_id()}:") - elif isinstance(event, core_events.DeviceUpdate): - _pending_attr_changes.add("#:") - else: - from qtoggleserver.slaves import events as slaves_events - - if isinstance( - event, - (slaves_events.SlaveDeviceAdd, slaves_events.SlaveDeviceRemove, slaves_events.SlaveDeviceUpdate), - ): - _pending_attr_changes.add(f"#{event.get_slave().get_name()}:") +_attr_change_handler = main_utils.AttrChangeHandler() def _get_changed_time_deps(now_int: int) -> tuple[bool, set[str]]: @@ -136,17 +117,17 @@ async def read_ports(ports_to_read: list[core_ports.BasePort] | None = None) -> if not ports_to_read: second_changed, changes = _get_changed_time_deps(now_int) + + # Get (and clear) any attribute dep changes (port or device) flagged by the event handler since last call + pending_attr_changes = _attr_change_handler.pop_pending() + if pending_attr_changes: + changes.update(pending_attr_changes) else: # When `ports_to_read` are given, this call is made from outside of the main update loop. Don't touch # time-related deps unless called from main update loop. second_changed = False changes = {DEP_ASAP} - # Drain any attribute dep changes (port or device) flagged by the event handler since the last tick - if _pending_attr_changes: - changes.update(_pending_attr_changes) - _pending_attr_changes.clear() - all_ports = list(core_ports.get_all()) if not ports_to_read: ports_to_read = all_ports @@ -305,7 +286,7 @@ async def init() -> None: loop = asyncio.get_running_loop() - core_events.register_handler(_AttrChangeHandler()) + core_events.register_handler(_attr_change_handler) force_eval_expressions() _update_loop_task = loop.create_task(update_loop()) diff --git a/qtoggleserver/utils/main.py b/qtoggleserver/utils/main.py new file mode 100644 index 00000000..ddd5f550 --- /dev/null +++ b/qtoggleserver/utils/main.py @@ -0,0 +1,39 @@ +from qtoggleserver.core import events as core_events +from qtoggleserver.slaves import events as slaves_events + + +class AttrChangeHandler(core_events.Handler): + """Accumulate attribute-related dep strings triggered by port/device change events. + + Listens for structural changes (port/device add, remove, update) and records which expression dependency prefixes + may be stale. The accumulated set is consumed via `pop_pending` on each evaluation tick so that expressions + depending on those attrs are re-evaluated. + + Dep strings produced: + - ``$port_id:`` — for "port-add", "port-remove", "port-update" + - ``#:`` — for "device-update" + - ``#name:`` — for "slave-device-add", "slave-device-remove", "slave-device-update" + """ + + FIRE_AND_FORGET = False + + def __init__(self) -> None: + super().__init__() + self._pending: set[str] = set() + + def pop_pending(self) -> set[str]: + """Return the pending changes and clear the internal set.""" + pending = self._pending.copy() + self._pending.clear() + return pending + + 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()}:") + elif isinstance(event, core_events.DeviceUpdate): + self._pending.add("#:") + elif isinstance( + event, + (slaves_events.SlaveDeviceAdd, slaves_events.SlaveDeviceRemove, slaves_events.SlaveDeviceUpdate), + ): + self._pending.add(f"#{event.get_slave().get_name()}:") diff --git a/tests/unit/qtoggleserver/core/test_main.py b/tests/unit/qtoggleserver/core/test_main.py index 105f43d5..8466181a 100644 --- a/tests/unit/qtoggleserver/core/test_main.py +++ b/tests/unit/qtoggleserver/core/test_main.py @@ -512,61 +512,54 @@ 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.""" - core_main._pending_attr_changes.clear() + """Clear pending attr changes before and after each test.""" + core_main._attr_change_handler._pending.clear() yield - core_main._pending_attr_changes.clear() + core_main._attr_change_handler._pending.clear() async def test_port_update_adds_attr_dep(self, mock_num_port1): """PortUpdate event should add `$port_id:` to _pending_attr_changes.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(core_events.PortUpdate(mock_num_port1)) - assert core_main._pending_attr_changes == {"$nid1:"} + await core_main._attr_change_handler.handle_event(core_events.PortUpdate(mock_num_port1)) + assert core_main._attr_change_handler._pending == {"$nid1:"} async def test_port_add_adds_attr_dep(self, mock_num_port1): """PortAdd event should add `$port_id:` to _pending_attr_changes.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(core_events.PortAdd(mock_num_port1)) - assert core_main._pending_attr_changes == {"$nid1:"} + await core_main._attr_change_handler.handle_event(core_events.PortAdd(mock_num_port1)) + 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.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(core_events.PortRemove(mock_num_port1)) - assert core_main._pending_attr_changes == {"$nid1:"} + await core_main._attr_change_handler.handle_event(core_events.PortRemove(mock_num_port1)) + assert core_main._attr_change_handler._pending == {"$nid1:"} async def test_device_update_adds_device_dep(self): """DeviceUpdate event should add `#:` to _pending_attr_changes.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(core_events.DeviceUpdate()) - assert core_main._pending_attr_changes == {"#:"} + await core_main._attr_change_handler.handle_event(core_events.DeviceUpdate()) + 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.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(core_events.ValueChange(None, 42, mock_num_port1)) - assert core_main._pending_attr_changes == set() + await core_main._attr_change_handler.handle_event(core_events.ValueChange(None, 42, mock_num_port1)) + 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.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(core_events.PortUpdate(mock_num_port1)) - await handler.handle_event(core_events.PortUpdate(mock_num_port2)) - assert core_main._pending_attr_changes == {"$nid1:", "$nid2:"} + 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:"} async def test_port_and_device_events_accumulate(self, mock_num_port1): """Port and device events should both accumulate into _pending_attr_changes.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(core_events.PortUpdate(mock_num_port1)) - await handler.handle_event(core_events.DeviceUpdate()) - assert core_main._pending_attr_changes == {"$nid1:", "#:"} + 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:", "#:"} 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.""" @@ -574,7 +567,7 @@ async def test_read_ports_drains_pending_attr_changes(self, freezer, mocker, moc freezer.move_to(dummy_utc_datetime) await read_ports() # prime time state - core_main._pending_attr_changes.add("$nid1:") + core_main._attr_change_handler._pending.add("$nid1:") spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() @@ -587,7 +580,7 @@ async def test_read_ports_drains_device_dep(self, freezer, mocker, mock_num_port freezer.move_to(dummy_utc_datetime) await read_ports() # prime time state - core_main._pending_attr_changes.add("#:") + core_main._attr_change_handler._pending.add("#:") spy = mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() @@ -595,16 +588,16 @@ async def test_read_ports_drains_device_dep(self, freezer, mocker, mock_num_port assert "#:" in call_changes async def test_read_ports_clears_pending_after_drain(self, freezer, mocker, mock_num_port1, dummy_utc_datetime): - """_pending_attr_changes must be empty after read_ports() drains it.""" + """_pending must be empty after read_ports() drains it.""" freezer.move_to(dummy_utc_datetime) await read_ports() # prime time state - core_main._pending_attr_changes.add("$nid1:") + core_main._attr_change_handler._pending.add("$nid1:") mocker.patch("qtoggleserver.core.main._eval_changed_expressions") await read_ports() - assert core_main._pending_attr_changes == set() + assert core_main._attr_change_handler._pending == set() async def test_port_attr_dep_triggers_expression_eval(self, mocker, mock_num_port1, mock_num_port2): """A port whose expression references $nid1: should be evaluated when `$nid1:` is in changes.""" @@ -631,31 +624,27 @@ 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.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1"))) - assert core_main._pending_attr_changes == {"#slave1:"} + await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1"))) + 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.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(slaves_events.SlaveDeviceAdd(_make_mock_slave("slave1"))) - assert core_main._pending_attr_changes == {"#slave1:"} + await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceAdd(_make_mock_slave("slave1"))) + 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.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(slaves_events.SlaveDeviceRemove(_make_mock_slave("slave1"))) - assert core_main._pending_attr_changes == {"#slave1:"} + await core_main._attr_change_handler.handle_event(slaves_events.SlaveDeviceRemove(_make_mock_slave("slave1"))) + 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.""" - handler = core_main._AttrChangeHandler() - await handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave1"))) - await handler.handle_event(slaves_events.SlaveDeviceUpdate(_make_mock_slave("slave2"))) - assert core_main._pending_attr_changes == {"#slave1:", "#slave2:"} + 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"))) + assert core_main._attr_change_handler._pending == {"#slave1:", "#slave2:"} async def test_slave_dep_triggers_expression_eval(self, mocker, mock_num_port1): """A port whose expression references #slave1:attr should be evaluated when `#slave1:` is in changes."""