From 05bc8bdf25d5443d2cbfd4596a77b0fd37c3afef Mon Sep 17 00:00:00 2001 From: Hector van der Aa Date: Wed, 8 Jul 2026 12:30:47 +0100 Subject: [PATCH] Working data ingress into live data buffer --- pyproject.toml | 2 +- src/dynalab/app.py | 24 ++++++- src/dynalab/data/json.py | 4 ++ src/dynalab/state.py | 3 + src/dynalab/ui/windows.py | 8 +-- test/minimal_connector.py | 129 +++++++++++++++++++++++++++----------- uv.lock | 16 +---- 7 files changed, 131 insertions(+), 55 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 31f39d9..0866cae 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -17,8 +17,8 @@ dynalab = "dynalab.app:run" package = true [tool.uv.sources] -dpg-gauges = { path = "../../dpg-gauges", editable = true } dynalab-protocol = { path = "../dynalab-protocol", editable = true } +dpg-gauges = { git = "https://git.h3cx.dev/h3cx/dpg-gauges.git" } [build-system] requires = ["setuptools>=68"] diff --git a/src/dynalab/app.py b/src/dynalab/app.py index f3eb86b..bc0cf5c 100644 --- a/src/dynalab/app.py +++ b/src/dynalab/app.py @@ -18,6 +18,7 @@ import dynalab.ui.windows from dynalab.ui.windows import build_loading, build_windows import dpg_gauges as dpgg +from dpg_gauges.exceptions import GaugeNotFoundError def _setup_logging(level: int = logging.INFO) -> None: @@ -168,10 +169,15 @@ def run() -> None: except Empty: pass else: - dpgg.analog_gauge( + dpgg.digital_gauge( tag=f"{new_signal.id}", label=f"{new_signal.name}", parent=MENU_WINDOW_DATA_INPUT, + unit=new_signal.unit if new_signal.unit is not None else "", + value=0, + precision=1, + min_value=new_signal.min_value, + max_value=new_signal.max_value, ) try: @@ -181,6 +187,22 @@ def run() -> None: else: dpgg.delete_gauge(f"{removed_signal.id}") + ctr: int = 0 + while not state.new_data_id.empty() or ctr < 5: + try: + new_value_id = state.new_data_id.get_nowait() + except Empty: + pass + else: + try: + dpgg.set_value( + tag=f"{new_value_id}", value=state.live_data[new_value_id] + ) + except GaugeNotFoundError: + pass + + ctr += 1 + # Render frame dpg.render_dearpygui_frame() finally: diff --git a/src/dynalab/data/json.py b/src/dynalab/data/json.py index c78de47..90a03df 100644 --- a/src/dynalab/data/json.py +++ b/src/dynalab/data/json.py @@ -13,6 +13,7 @@ from dynalab_protocol.models import ( HandshakeAccepted, HandshakeRejected, Heartbeat, + ValueDescriptor, VersionDescriptor, machine_timestamp_ms, ) @@ -100,6 +101,9 @@ async def _handle_connector( state.last_heartbeat[connector_hello.connector_uuid] = ( message.return_timestamp ) + if isinstance(message, ValueDescriptor): + state.live_data[message.signal_id] = message.value + state.new_data_id.put_nowait(message.signal_id) except TimeoutError: continue except ConnectionError: diff --git a/src/dynalab/state.py b/src/dynalab/state.py index 5eb5ba0..4866fb0 100644 --- a/src/dynalab/state.py +++ b/src/dynalab/state.py @@ -16,3 +16,6 @@ class AppState: last_heartbeat: dict[UUID, int] = field(default_factory=dict) new_signal_queue: Queue[SignalDescriptor] = Queue() removed_signal_queue: Queue[SignalDescriptor] = Queue() + + live_data: dict[UUID, float] = field(default_factory=dict) + new_data_id: Queue[UUID] = Queue() diff --git a/src/dynalab/ui/windows.py b/src/dynalab/ui/windows.py index b3d8e63..97830e4 100644 --- a/src/dynalab/ui/windows.py +++ b/src/dynalab/ui/windows.py @@ -1,6 +1,7 @@ # Copyright (C) 2026 Hector van der Aa # Copyright (C) 2026 Association Exergie # SPDX-License-Identifier: GPL-3.0-or-later +from os import wait import dearpygui.dearpygui as dpg import dpg_gauges as dpgg @@ -27,12 +28,9 @@ def build_windows() -> None: ): # with dpg.group(tag=MENU_WINDOW_DATA_INPUT): with dpgg.gauge_panel( - label="Autobreak narrow panel", width=430, height=285 + label="Autobreak narrow panel", autosize_x=True, autosize_y=True ): with dpgg.gauge_grid( columns="auto", min_column_width=180, tag=MENU_WINDOW_DATA_INPUT ): - dpgg.analog_gauge( - tag="test_gauge", - label="new_gauge", - ) + pass diff --git a/test/minimal_connector.py b/test/minimal_connector.py index dd62e31..c5dd302 100644 --- a/test/minimal_connector.py +++ b/test/minimal_connector.py @@ -2,6 +2,8 @@ from __future__ import annotations import asyncio +from random import random +from time import monotonic import uuid from dynalab_protocol.models import ( ConnectorHello, @@ -10,6 +12,7 @@ from dynalab_protocol.models import ( HandshakeRejected, Heartbeat, SignalDescriptor, + ValueDescriptor, machine_timestamp_ms, ) from dynalab_protocol.wire import read_message, write_message @@ -17,9 +20,79 @@ from dynalab_protocol.wire import read_message, write_message HOST = "127.0.0.1" PORT = 8765 +RETRY_DELAY_SECONDS = 2 -async def main() -> None: +def build_signals() -> list[SignalDescriptor]: + """Describe the signals this example connector exposes to DynaLab.""" + return [ + SignalDescriptor( + id=uuid.uuid4(), name="lambda", type="number", min_value=0, max_value=5 + ), + SignalDescriptor( + id=uuid.uuid4(), + name="torque", + type="number", + min_value=0, + max_value=10, + unit="Nm", + ), + SignalDescriptor( + id=uuid.uuid4(), + name="RPM", + type="number", + min_value=0, + max_value=6000, + unit="RPM", + ), + SignalDescriptor( + id=uuid.uuid4(), + name="temperature", + type="number", + min_value=0, + max_value=120, + unit="°C", + ), + SignalDescriptor( + id=uuid.uuid4(), + name="voltage", + type="number", + min_value=0, + max_value=20, + unit="V", + ), + SignalDescriptor( + id=uuid.uuid4(), + name="current", + type="number", + min_value=0, + max_value=50, + unit="A", + ), + ] + + +async def _value_sender( + writer: asyncio.StreamWriter, connector_hello: ConnectorHello +) -> None: + last_value_send = monotonic() + while True: + if monotonic() > last_value_send + 0.1: + multipliers: list[int] = [5, 10, 6000, 120, 20, 50] + for i in range(6): + random_val = random() * multipliers[i] + value = ValueDescriptor( + signal_id=connector_hello.signals[i].id, value=random_val + ) + await write_message(writer, value) + last_value_send = monotonic() + + else: + await asyncio.sleep(0.05) + + +async def run_connector_session() -> None: + """Open one connection to DynaLab and run the connector protocol.""" reader, writer = await asyncio.open_connection(HOST, PORT) dynalab_hello: DynaLabHello @@ -29,6 +102,7 @@ async def main() -> None: print(f"Connected to DynaLab at {peer}") try: + # DynaLab starts the protocol by introducing itself to the connector. message = await asyncio.wait_for(read_message(reader), timeout=5.0) if not isinstance(message, DynaLabHello): @@ -37,47 +111,16 @@ async def main() -> None: dynalab_hello = message print(f"Received DynaLab Hello\n {dynalab_hello}") - signals: list[SignalDescriptor] = [ - SignalDescriptor( - id=uuid.uuid4(), name="lambda", type="number", min_value=0, max_value=5 - ), - SignalDescriptor( - id=uuid.uuid4(), name="torque", type="number", min_value=0, max_value=10 - ), - SignalDescriptor( - id=uuid.uuid4(), name="RPM", type="number", min_value=0, max_value=6000 - ), - SignalDescriptor( - id=uuid.uuid4(), - name="temperature", - type="number", - min_value=0, - max_value=120, - ), - SignalDescriptor( - id=uuid.uuid4(), - name="voltage", - type="number", - min_value=0, - max_value=20, - ), - SignalDescriptor( - id=uuid.uuid4(), - name="current", - type="number", - min_value=0, - max_value=50, - ), - ] - + # Reply with the connector identity and signal list. connector_hello = ConnectorHello( connector_uuid=uuid.uuid4(), connector_name="Minimal Connector", connector_version="a0.0.1", - signals=signals, + signals=build_signals(), ) await write_message(writer, connector_hello) + # DynaLab accepts or rejects the connector after validating the hello. message = await asyncio.wait_for(read_message(reader), timeout=5.0) if isinstance(message, HandshakeAccepted): @@ -87,6 +130,8 @@ async def main() -> None: print(f"Connector handshake rejected: {message.reason}") return + value_task = asyncio.create_task(_value_sender(writer, connector_hello)) + # Keep the connection alive by returning heartbeat messages. while True: message = await read_message(reader) print(f"Received {message}") @@ -100,11 +145,25 @@ async def main() -> None: raise finally: + value_task.cancel() print("Disconnecting") writer.close() await writer.wait_closed() +async def main() -> None: + while True: + try: + await run_connector_session() + return + except ConnectionError as error: + print( + f"Connection error: {error}. " + f"Retrying in {RETRY_DELAY_SECONDS} seconds..." + ) + await asyncio.sleep(RETRY_DELAY_SECONDS) + + if __name__ == "__main__": try: asyncio.run(main()) diff --git a/uv.lock b/uv.lock index 14d1b57..4233189 100644 --- a/uv.lock +++ b/uv.lock @@ -28,22 +28,12 @@ wheels = [ [[package]] name = "dpg-gauges" -version = "1.1.0" -source = { editable = "../../dpg-gauges" } +version = "1.1.2" +source = { git = "https://git.h3cx.dev/h3cx/dpg-gauges.git#72b6e3714fce707617b2c9273aac7105db5eed74" } dependencies = [ { name = "dearpygui" }, ] -[package.metadata] -requires-dist = [{ name = "dearpygui", specifier = ">=2.3.1" }] - -[package.metadata.requires-dev] -dev = [ - { name = "pyright", specifier = ">=1.1.411" }, - { name = "pytest", specifier = ">=9.1.1" }, - { name = "ruff", specifier = ">=0.15.20" }, -] - [[package]] name = "dynalab" version = "0.1.0" @@ -57,7 +47,7 @@ dependencies = [ [package.metadata] requires-dist = [ { name = "dearpygui", specifier = ">=2.3.1" }, - { name = "dpg-gauges", editable = "../../dpg-gauges" }, + { name = "dpg-gauges", git = "https://git.h3cx.dev/h3cx/dpg-gauges.git" }, { name = "dynalab-protocol", editable = "../dynalab-protocol" }, ]