|
6 | 6 | from qtoggleserver.core import expressions as core_expressions |
7 | 7 | from qtoggleserver.core import ports as core_ports |
8 | 8 | from qtoggleserver.core.expressions import exceptions as expression_exceptions |
| 9 | +from qtoggleserver.core.expressions.exceptions import ExpressionEvalException |
9 | 10 | from qtoggleserver.core.typing import Attribute, Attributes, NullablePortValue |
10 | 11 | from qtoggleserver.slaves import devices as slaves_devices |
11 | 12 | from qtoggleserver.slaves import events as slaves_events |
@@ -41,6 +42,8 @@ def __init__(self, *, filter: dict | None = None, name: str | None = None) -> No |
41 | 42 | self._filter_slave_attr_transitions: dict[str, tuple[Attribute, Attribute]] = {} |
42 | 43 | self._filter_slave_attr_names: set[str] = set() |
43 | 44 |
|
| 45 | + self._filter_expression: core_expressions.Expression | None = None |
| 46 | + |
44 | 47 | # Maintain an internal "last" state for all objects, so we can detect changes in attributes and values |
45 | 48 | self._device_attrs: Attributes | dict[str, list[Attribute]] = {} |
46 | 49 | self._port_values: dict[str, NullablePortValue] = {} |
@@ -113,6 +116,18 @@ def _prepare_filter(self) -> None: |
113 | 116 | self._filter_slave_attr_names.update(self._filter_slave_attrs.keys()) |
114 | 117 | self._filter_slave_attr_names.update(self._filter_slave_attr_transitions.keys()) |
115 | 118 |
|
| 119 | + filter_sexpression = self._filter.get("expression") |
| 120 | + if isinstance(filter_sexpression, str): |
| 121 | + try: |
| 122 | + self.debug('using filter expression "%s"', filter_sexpression) |
| 123 | + self._filter_expression = core_expressions.parse( |
| 124 | + self_port_id=None, sexpression=filter_sexpression, role=core_expressions.Role.FILTER |
| 125 | + ) |
| 126 | + except expression_exceptions.ExpressionParseError as e: |
| 127 | + self.error('failed to parse filter expression "%s": %s', filter_sexpression, e) |
| 128 | + |
| 129 | + raise |
| 130 | + |
116 | 131 | self._filter_prepared = True |
117 | 132 | self.debug("filter prepared") |
118 | 133 |
|
@@ -261,7 +276,11 @@ async def accepts_port_value( |
261 | 276 | elif isinstance(self._filter_port_value, core_expressions.Expression): # an expression |
262 | 277 | port_values = {p.get_id(): p.get_last_read_value() for p in core_ports.get_all() if p.is_enabled()} |
263 | 278 | eval_context = core_expressions.EvalContext(port_values=port_values, now_ms=int(time.time() * 1000)) |
264 | | - if new_value != await self._filter_port_value.eval(context=eval_context): |
| 279 | + try: |
| 280 | + if new_value != await self._filter_port_value.eval(context=eval_context): |
| 281 | + return False |
| 282 | + except ExpressionEvalException as e: |
| 283 | + self.warning('Expression evaluation failed for "%s": %s', self._filter_port_value, e) |
265 | 284 | return False |
266 | 285 | elif new_value != self._filter_port_value: |
267 | 286 | return False |
@@ -322,6 +341,16 @@ async def accepts( |
322 | 341 | ): |
323 | 342 | return False |
324 | 343 |
|
| 344 | + if self._filter_expression: |
| 345 | + port_values = {p.get_id(): p.get_last_read_value() for p in core_ports.get_all() if p.is_enabled()} |
| 346 | + eval_context = core_expressions.EvalContext(port_values=port_values, now_ms=int(time.time() * 1000)) |
| 347 | + try: |
| 348 | + if not await self._filter_expression.eval(context=eval_context): |
| 349 | + return False |
| 350 | + except ExpressionEvalException as e: |
| 351 | + self.warning('Expression evaluation failed for "%s": %s', self._filter_expression, e) |
| 352 | + return False |
| 353 | + |
325 | 354 | return True |
326 | 355 |
|
327 | 356 | def get_device_attrs(self) -> Attributes: |
|
0 commit comments