Files
2026-09-11 23:45:27 +02:00

109 lines
3.7 KiB
Python

import logging
from queue import Empty, Queue
import time
from uuid import UUID, uuid4
import pytest
from dynalab_core.derive import DeriveRegistry
from dynalab_core.protocols.packets.data import ValueDescriptor
from dynalab_core.protocols.packets.handshake import SignalDescriptor
PROCESSING_FUNCTION = """
from dynalab_core.protocols.packets.data import ValueDescriptor
from dynalab_core.protocols.packets.handshake import SignalDescriptor
def add(left: ValueDescriptor, right: ValueDescriptor, output: SignalDescriptor) -> ValueDescriptor:
return ValueDescriptor(
signal_id=output.id,
value=left.value + right.value,
timestamp=left.timestamp,
)
"""
FAILING_PROCESSING_FUNCTION = """
from dynalab_core.protocols.packets.data import ValueDescriptor
from dynalab_core.protocols.packets.handshake import SignalDescriptor
def fail(value: ValueDescriptor, output: SignalDescriptor) -> ValueDescriptor:
raise RuntimeError("derive failed")
"""
def test_derive_registry_routes_with_live_values() -> None:
output_queue: Queue[ValueDescriptor] = Queue()
left_signal = SignalDescriptor(id=uuid4(), name="Left", type="number")
right_signal = SignalDescriptor(id=uuid4(), name="Right", type="number")
output_signal = SignalDescriptor(id=uuid4(), name="Total", type="number")
live_values: dict[UUID, ValueDescriptor] = {}
registry = DeriveRegistry(output_queue, live_values.get)
try:
unit = registry.register(
PROCESSING_FUNCTION, [left_signal, right_signal], output_signal
)
right_value = ValueDescriptor(signal_id=right_signal.id, value=2.0, timestamp=1)
live_values[right_signal.id] = right_value
left_value = ValueDescriptor(signal_id=left_signal.id, value=3.0, timestamp=2)
registry.put_data(left_value)
deadline = time.monotonic() + 1.0
while True:
try:
result = output_queue.get_nowait()
break
except Empty:
if time.monotonic() >= deadline:
raise AssertionError("Derive unit did not produce a value")
time.sleep(0.01)
registry.unregister(unit.uuid())
assert registry.get(unit.uuid()) is None
finally:
registry.stop()
assert result.signal_id == output_signal.id
assert result.value == 5.0
assert result.timestamp == left_value.timestamp
def test_derive_registry_ignores_unknown_unit_on_unregister() -> None:
registry = DeriveRegistry(Queue(), lambda signal_id: None)
try:
registry.unregister(uuid4())
finally:
registry.stop()
def test_derive_worker_logs_failure_and_stops(
caplog: pytest.LogCaptureFixture,
) -> None:
input_signal = SignalDescriptor(id=uuid4(), name="Input", type="number")
output_signal = SignalDescriptor(id=uuid4(), name="Output", type="number")
registry = DeriveRegistry(Queue(), lambda signal_id: None)
with caplog.at_level(logging.DEBUG, logger="dynalab_core"):
try:
unit = registry.register(
FAILING_PROCESSING_FUNCTION, [input_signal], output_signal
)
unit.put_data(
[ValueDescriptor(signal_id=input_signal.id, value=1.0, timestamp=1)]
)
assert unit._stopped_event.wait(1.0)
assert registry.get_all_signal_descriptors() == []
finally:
registry.stop()
failures = [
record
for record in caplog.records
if getattr(record, "event", None) == "derive.worker_failed"
]
assert len(failures) == 1
assert failures[0].derive_uuid == str(unit.uuid())
assert failures[0].exception_type == "RuntimeError"
assert failures[0].exc_info is not None