Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ package = true

[tool.ruff]
line-length = 120
lint.extend-select = ["I", "RUF022", "ANN"]
lint.extend-select = ["I", "RUF022", "ANN", "E501"]
lint.extend-ignore = ["ANN002", "ANN003", "ANN401"]
lint.isort.lines-after-imports = 2
lint.isort.lines-between-types = 1
Expand Down
1 change: 1 addition & 0 deletions qtoggleserver/core/api/funcs/various.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ async def get_listen(request: core_api.APIRequest) -> GenericJSONList:
events = await session.reset_and_wait(timeout, request.access_level)
except CancelledError:
session.debug("waiting cancelled")
session.cancel()
return []

return [await e.to_json() for e in events]
Expand Down
9 changes: 3 additions & 6 deletions qtoggleserver/core/expressions/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,12 +50,9 @@


def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int = 1) -> Expression:
while sexpression and sexpression[0].isspace():
sexpression = sexpression[1:]
pos += 1

while sexpression and sexpression[-1].isspace():
sexpression = sexpression[:-1]
stripped = sexpression.lstrip()
pos += len(sexpression) - len(stripped)
sexpression = stripped.rstrip()

if sexpression and sexpression[0] in ("$", "@"):
return PortExpression.parse(self_port_id, sexpression, role, pos)
Expand Down
10 changes: 2 additions & 8 deletions qtoggleserver/core/expressions/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,13 +75,7 @@ class EvalContext:
def __init__(self, port_values: dict[str, NullablePortValue], now_ms: int = 0) -> None:
self.port_values: dict[str, NullablePortValue] = port_values
self.now_ms: int = now_ms
self.timestamp: int = now_ms // 1000

@property
def timestamp(self) -> int:
return int(self.now_ms / 1000)

def __str__(self) -> str:
return f"EvalContext(now_ms={self.now_ms}, port_values={self.port_values})"


EvalResult: TypeAlias = bool | int | float | str
EvalResult: TypeAlias = int | float | str
Comment thread
ccrisan marked this conversation as resolved.
12 changes: 7 additions & 5 deletions qtoggleserver/core/expressions/functions.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ def __init__(self, args: list[Expression], role: Role) -> None:
super().__init__(role)

self.args: list[Expression] = args
self._function_args: list[Function] = [arg for arg in args if isinstance(arg, Function)]

def __str__(self) -> str:
s = getattr(self, "_str", None)
Expand All @@ -56,11 +57,12 @@ def is_asap_eval_paused(self, now_ms: int) -> bool:
return now_ms < self._get_min_asap_eval_paused_until_ms()

def _get_min_asap_eval_paused_until_ms(self) -> int:
min_asap_eval_paused_until_ms = self._asap_eval_paused_until_ms if DEP_ASAP in self.DEPS else 1e13
min_asap_eval_paused_until_ms_args = [
arg._get_min_asap_eval_paused_until_ms() for arg in self.args if isinstance(arg, Function)
]
return min([min_asap_eval_paused_until_ms, *min_asap_eval_paused_until_ms_args])
result = self._asap_eval_paused_until_ms if DEP_ASAP in self.DEPS else int(1e13)
for arg in self._function_args:
child = arg._get_min_asap_eval_paused_until_ms()
if child < result:
result = child
return result

async def eval_args(self, context: EvalContext) -> list[EvalResult]:
return list(await asyncio.gather(*(a.eval(context) for a in self.args)))
Expand Down
25 changes: 10 additions & 15 deletions qtoggleserver/core/expressions/literalvalues.py
Original file line number Diff line number Diff line change
@@ -1,38 +1,33 @@
import re

from qtoggleserver.core.typing import NullablePortValue as CoreNullablePortValue
from qtoggleserver.core.typing import NullablePortValue

from .base import EvalContext, EvalResult, Expression, Role
from .exceptions import EmptyExpression, UnexpectedCharacter, ValueUnavailable


class LiteralValue(Expression):
def __init__(self, value: CoreNullablePortValue, sexpression: str, role: Role) -> None:
def __init__(self, value: NullablePortValue, sexpression: str, role: Role) -> None:
super().__init__(role)

self.value: CoreNullablePortValue = value
self.sexpression: str = sexpression
self._coerced_value: EvalResult | None = value
if isinstance(value, bool):
self._coerced_value = int(value)

def __str__(self) -> str:
return self.sexpression

async def _eval(self, context: EvalContext) -> EvalResult:
if self.value is None:
if self._coerced_value is None:
raise ValueUnavailable

if isinstance(self.value, int):
return self.value
else:
return float(self.value)
return self._coerced_value

@staticmethod
def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> Expression:
while sexpression and sexpression[0].isspace():
sexpression = sexpression[1:]
pos += 1

while sexpression and sexpression[-1].isspace():
sexpression = sexpression[:-1]
stripped = sexpression.lstrip()
pos += len(sexpression) - len(stripped)
sexpression = stripped.rstrip()

if not sexpression:
raise EmptyExpression()
Expand Down
20 changes: 10 additions & 10 deletions qtoggleserver/core/expressions/ports.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,20 +16,20 @@ def __init__(self, port_id: str, prefix: str, role: Role) -> None:

self.port_id: str = port_id
self.prefix: str = prefix
self._cached_port: core_ports.BasePort | None = None

def get_port(self) -> core_ports.BasePort:
return core_ports.get(self.port_id)
def get_port(self) -> core_ports.BasePort | None:
port = self._cached_port
if port is None or port.is_removed():
port = core_ports.get(self.port_id)
self._cached_port = port
return port

@staticmethod
def parse(self_port_id: str | None, sexpression: str, role: Role, pos: int) -> Expression:
# Remove leading whitespace
while sexpression and sexpression[0].isspace():
sexpression = sexpression[1:]
pos += 1

# Remove trailing whitespace
while sexpression and sexpression[-1].isspace():
sexpression = sexpression[:-1]
stripped = sexpression.lstrip()
pos += len(sexpression) - len(stripped)
sexpression = stripped.rstrip()

prefix = sexpression[0]
port_id = sexpression[1:]
Expand Down
32 changes: 13 additions & 19 deletions qtoggleserver/core/expressions/timeprocessing.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
from collections import deque

from . import TIME_JUMP_THRESHOLD
from .base import DEP_ASAP, EvalContext, EvalResult
from .functions import Function, function
Expand All @@ -13,7 +15,7 @@ class DelayFunction(Function):
def __init__(self, *args, **kwargs) -> None:
super().__init__(*args, **kwargs)

self._queue: list[tuple[int, float]] = []
self._queue: deque[tuple[int, float]] = deque(maxlen=self.HISTORY_SIZE)
self._last_value: float | None = None
self._current_value: float | None = None

Expand All @@ -28,16 +30,11 @@ async def _eval(self, context: EvalContext) -> EvalResult:
# Detect value transitions and build history
if value != self._last_value:
self._last_value = value

# Drop elements from queue if history size reached
while len(self._queue) >= self.HISTORY_SIZE:
self._queue.pop(0)

self._queue.append((context.now_ms, value))

# Process history
while self._queue and (context.now_ms - self._queue[0][0]) >= delay:
self._current_value = self._queue.pop(0)[1]
self._current_value = self._queue.popleft()[1]

if self._queue:
self.pause_asap_eval(self._queue[0][0] + delay)
Expand Down Expand Up @@ -169,7 +166,7 @@ async def _eval(self, context: EvalContext) -> EvalResult:
self._start_time_ms = 0 # stop timer
self.pause_asap_eval()

return self._start_time_ms > 0 and context.now_ms - self._start_time_ms >= duration
return int(self._start_time_ms > 0 and context.now_ms - self._start_time_ms >= duration)


@function("DERIV")
Expand Down Expand Up @@ -256,13 +253,13 @@ class FMAvgFunction(Function):
def __init__(self, *args, **kwargs) -> None:
super().__init__(*args, **kwargs)

self._queue: list[float] = []
self._queue: deque[float] = deque()
self._last_time_ms: int = 0
self._last_result: float = 0

async def _eval(self, context: EvalContext) -> EvalResult:
value, width, sampling_interval = await self.eval_args(context)
width = min(width, self.MAX_QUEUE_SIZE)
width = max(1, min(int(width), self.MAX_QUEUE_SIZE))

Comment thread
ccrisan marked this conversation as resolved.
if self._last_time_ms > 0:
if context.now_ms - self._last_time_ms < sampling_interval:
Expand All @@ -271,13 +268,12 @@ async def _eval(self, context: EvalContext) -> EvalResult:

# Make room for the new element
while len(self._queue) >= width:
self._queue.pop(0)
self._queue.popleft()

Comment thread
ccrisan marked this conversation as resolved.
self._queue.append(value)
self._last_time_ms = context.now_ms

queue = self._queue[-int(width) :]
self._last_result = sum(queue) / len(queue)
self._last_result = sum(self._queue) / len(self._queue)
self.pause_asap_eval(self._last_time_ms + sampling_interval)

return self._last_result
Expand All @@ -293,30 +289,28 @@ class FMedianFunction(Function):
def __init__(self, *args, **kwargs) -> None:
super().__init__(*args, **kwargs)

self._queue: list[float] = []
self._queue: deque[float] = deque()
self._last_time_ms: int = 0
self._last_result: float = 0

async def _eval(self, context: EvalContext) -> EvalResult:
value, width, sampling_interval = await self.eval_args(context)
width = min(width, self.MAX_QUEUE_SIZE)
width = max(1, min(int(width), self.MAX_QUEUE_SIZE))

Comment thread
ccrisan marked this conversation as resolved.
if self._last_time_ms > 0:
if context.now_ms - self._last_time_ms < sampling_interval:
self.pause_asap_eval(self._last_time_ms + sampling_interval)
self.pause_asap_eval(self._last_time_ms + sampling_interval)
return self._last_result

# Make room for the new element
while len(self._queue) >= width:
self._queue.pop(0)
self._queue.popleft()

Comment thread
ccrisan marked this conversation as resolved.
self._queue.append(value)
self._last_time_ms = context.now_ms
self.pause_asap_eval(self._last_time_ms + sampling_interval)

queue = self._queue[-int(width) :]
queue.sort()
queue = sorted(self._queue)
self._last_result = queue[len(queue) // 2]

return self._last_result
38 changes: 15 additions & 23 deletions qtoggleserver/core/expressions/various.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,9 @@ class AvailableFunction(Function):

async def _eval(self, context: EvalContext) -> EvalResult:
try:
return await self.args[0].eval(context) is not None
return int(await self.args[0].eval(context) is not None)
except ValueUnavailable:
return False
return 0


@function("DEFAULT")
Expand Down Expand Up @@ -58,9 +58,9 @@ def __init__(self, *args, **kwargs) -> None:
async def _eval(self, context: EvalContext) -> EvalResult:
value = await self.args[0].eval(context)

result = False
result = 0
if self._last_value is not None and value > self._last_value:
result = True
result = 1
self._last_value = value

return result
Expand All @@ -78,9 +78,9 @@ def __init__(self, *args, **kwargs) -> None:
async def _eval(self, context: EvalContext) -> EvalResult:
value = await self.args[0].eval(context)

result = False
result = 0
if self._last_value is not None and value < self._last_value:
result = True
result = 1
self._last_value = value

return result
Expand Down Expand Up @@ -154,9 +154,9 @@ class OnOffAutoFunction(Function):
async def _eval(self, context: EvalContext) -> EvalResult:
value, auto = await self.eval_args(context)
if value > 0:
return True
return 1
elif value < 0:
return False
return 0
else:
return auto

Expand All @@ -171,33 +171,25 @@ def __init__(self, *args, **kwargs) -> None:
super().__init__(*args, **kwargs)

self._start_time_ms: int = 0
self._num_values: int = len(self.args) // 2

async def _eval(self, context: EvalContext) -> EvalResult:
if self._start_time_ms == 0:
self._start_time_ms = context.now_ms

args = await self.eval_args(context)
num_values = len(args) // 2
values = []
delays = []
total_delay = 0
for i in range(num_values):
values.append(args[i * 2])
delays.append(args[i * 2 + 1])
for i in range(self._num_values):
total_delay += args[i * 2 + 1]

if len(delays) < len(values):
delays.append(0)

delta = context.now_ms - self._start_time_ms
delta = delta % total_delay # work modulo total_delay, to create repeat effect
delta = (context.now_ms - self._start_time_ms) % total_delay
delay_so_far = 0
result = values[0]
for i in range(num_values):
delay_so_far += delays[i]
result = args[0]
for i in range(self._num_values):
delay_so_far += args[i * 2 + 1]
if delay_so_far >= delta:
self.pause_asap_eval(context.now_ms + delay_so_far - delta)
result = values[i]
result = args[i * 2]
break

return result
Expand Down
8 changes: 5 additions & 3 deletions qtoggleserver/core/ports.py
Original file line number Diff line number Diff line change
Expand Up @@ -256,6 +256,7 @@ def __init__(self, port_id: str) -> None:
self._pending_save: bool = False

self._loaded: bool = False
self._removed: bool = False
self._after_set_attr_debounced = Debounced(self._after_set_attr)

def __str__(self) -> str:
Expand Down Expand Up @@ -879,7 +880,7 @@ async def load(self) -> None:
data = await persist.get(self.PERSIST_COLLECTION, self.get_id()) or {}
await self.load_from_data(data)

self.set_loaded()
self._loaded = True
self.initialize()

async def reset(self) -> None:
Expand Down Expand Up @@ -1012,13 +1013,14 @@ async def cleanup(self) -> None:
def is_loaded(self) -> bool:
return self._loaded

def set_loaded(self) -> None:
self._loaded = True
def is_removed(self) -> bool:
return self._removed

async def remove(self, persisted_data: bool = True) -> None:
await self.cleanup()

self.debug("removing port")
self._removed = True
_ports_by_id.pop(self._id, None)

if persisted_data:
Expand Down
Loading
Loading