From ebe7e92e2238e414ebfae1c00cbc2314a169890b Mon Sep 17 00:00:00 2001 From: Hector van der Aa Date: Sat, 22 Aug 2026 17:45:35 +0100 Subject: [PATCH] Added origin descriptor to SignalDescriptor and adapted connector registry for loopback connector redesign SignalDescriptor now have a new origin field to distinguish "source" from "derived" signals, this is ahead of the internal loopback system for data processing, all external conectors must use source, derived is targeted only for internal use The connector registry has been adapted to now have a unique internal loopback, this does not have an associated endpoint so it is stored apart and all relevant functions have been updated to include it for proper handling of signals by DLPak etc --- src/dynalab_core/__init__.py | 4 ++++ src/dynalab_core/buffer.py | 9 +++++++ src/dynalab_core/constants.py | 10 +++++++- src/dynalab_core/protocols/endpoint.py | 24 +++++++++++++++++-- .../protocols/packets/handshake.py | 1 + test/manual/peer.py | 4 +++- 6 files changed, 48 insertions(+), 4 deletions(-) diff --git a/src/dynalab_core/__init__.py b/src/dynalab_core/__init__.py index 3a08b95..18528e1 100644 --- a/src/dynalab_core/__init__.py +++ b/src/dynalab_core/__init__.py @@ -8,6 +8,7 @@ import logging from queue import Empty, Queue import threading from threading import Lock, Thread +import time from typing import Literal from uuid import UUID @@ -116,12 +117,15 @@ class Core: def stop_recording(self) -> None: self._recording.clear() + time.sleep(1) + self._recording_buffer.normalize() self._processing_buffer = DLPak() self._processing_buffer.set_data(self._recording_buffer) self._processing_buffer.set_manifest( self._recording_timestamp, self._connector_registry.get_all_signal_descriptors(), ) + # TODO: remove debug behavior default saving to output self._processing_buffer.write("./", "output") def get_signal_descriptor(self, signal_id: UUID) -> SignalDescriptor | None: diff --git a/src/dynalab_core/buffer.py b/src/dynalab_core/buffer.py index 0aec51f..2267f18 100644 --- a/src/dynalab_core/buffer.py +++ b/src/dynalab_core/buffer.py @@ -31,6 +31,15 @@ class ValueBuffer: with self._lock: return len(self._samples) + def normalize(self) -> None: + with self._lock: + minimum = min(sample[0] for sample in self._samples) + + self._samples = [ + (timestamp - minimum, signal_id, value) + for timestamp, signal_id, value in self._samples + ] + def export_csv(self) -> str: with self._lock: samples = list(self._samples) diff --git a/src/dynalab_core/constants.py b/src/dynalab_core/constants.py index 9893b29..8506ffb 100644 --- a/src/dynalab_core/constants.py +++ b/src/dynalab_core/constants.py @@ -5,7 +5,7 @@ from uuid import uuid4 from dynalab_core.protocols.common import VersionDescriptor from dynalab_core.protocols.constants import PROTOCOL_VERSION -from dynalab_core.protocols.packets.handshake import DynaLabHello +from dynalab_core.protocols.packets.handshake import ConnectorHello, DynaLabHello CORE_VERSION = VersionDescriptor(type="alpha", major=0, minor=0, patch=1) @@ -17,3 +17,11 @@ HELLO_PACKET = DynaLabHello( heartbeat_interval_ms=1000, heartbeat_timeout_ms=5000, ) + +INTERNAL_CONNECTOR_HELLO = ConnectorHello( + connector_uuid=uuid4(), + protocol_version=PROTOCOL_VERSION, + connector_name="INTERNAL_LOOPBACK", + connector_version=CORE_VERSION.get_version(), + signals=[], +) diff --git a/src/dynalab_core/protocols/endpoint.py b/src/dynalab_core/protocols/endpoint.py index fd061b4..64fd843 100644 --- a/src/dynalab_core/protocols/endpoint.py +++ b/src/dynalab_core/protocols/endpoint.py @@ -9,7 +9,7 @@ from threading import RLock, Thread from time import monotonic, monotonic_ns, sleep from uuid import UUID -from dynalab_core.constants import HELLO_PACKET +from dynalab_core.constants import HELLO_PACKET, INTERNAL_CONNECTOR_HELLO from dynalab_core.protocols.errors import ( ConnectorEndpointQueueFullError, ConnectorRegistryAlreadyRegisteredError, @@ -322,9 +322,15 @@ class ConnectorEndpoint: class ConnectorRegistry: def __init__(self, core_input_queue: Queue) -> None: self._endpoints: dict[UUID, ConnectorEndpoint] = {} + self._internal_connector: ConnectorHello = INTERNAL_CONNECTOR_HELLO self._lock = RLock() self._core_input_queue = core_input_queue + def add_internal_signal(self, signal: SignalDescriptor) -> None: + with self._lock: + if not any(s.id == signal.id for s in self._internal_connector.signals): + self._internal_connector.signals.append(signal) + def register( self, hello: ConnectorHello, @@ -402,7 +408,14 @@ class ConnectorRegistry: def get(self, connector_uuid: UUID) -> ConnectorEndpoint | None: with self._lock: - return self._endpoints.get(connector_uuid) + endpoint = self._endpoints.get(connector_uuid) + if not endpoint: + return endpoint + + if self._internal_connector.connector_uuid == connector_uuid: + return self._internal_connector + + return None def get_signal_descriptor(self, signal_id: UUID) -> SignalDescriptor | None: with self._lock: @@ -413,6 +426,10 @@ class ConnectorRegistry: if signal is not None: return signal + for signal in self._internal_connector.signals: + if signal.id == signal_id: + return signal + return None def get_all_signal_descriptors(self) -> list[SignalDescriptor]: @@ -425,6 +442,9 @@ class ConnectorRegistry: for signal in endpoint_signals: signals.append(signal) + for signal in self._internal_connector.signals: + signals.append(signal) + return signals def stop(self) -> None: diff --git a/src/dynalab_core/protocols/packets/handshake.py b/src/dynalab_core/protocols/packets/handshake.py index 0ce2a0b..4e812e7 100644 --- a/src/dynalab_core/protocols/packets/handshake.py +++ b/src/dynalab_core/protocols/packets/handshake.py @@ -18,6 +18,7 @@ class SignalDescriptor(BaseModel): max_value: float | None = None unit: str | None = None timeout_ms: int = 2000 + origin: Literal["source", "derived"] = "source" class DynaLabHello(BaseModel): diff --git a/test/manual/peer.py b/test/manual/peer.py index cea0da0..9c90bd0 100644 --- a/test/manual/peer.py +++ b/test/manual/peer.py @@ -119,7 +119,9 @@ def value_sender_thread( message = ValueBatch(values=[]) for i in range(10): value = ValueDescriptor( - signal_id=signal_id, value=2.0, timestamp=monotonic_ns() + signal_id=signal_id, + value=2.0, + timestamp=monotonic_ns(), ) message.values.append(value)