From 1cd6bb695c89445020602da9b237ebaed115e2bf Mon Sep 17 00:00:00 2001 From: Hector van der Aa Date: Thu, 6 Aug 2026 15:29:02 +0100 Subject: [PATCH] Implemented initial DLPak design Implemented DLPak design with zip storage, DLPak consists of two files, a manifest.json which provides details on the signals, time of recording and currently blank comment field and data.csv which stores all the signals as csv with the nanosecond timestamp and indexed against the signal id --- src/dynalab_core/__init__.py | 14 +++++++- src/dynalab_core/buffer.py | 40 ++++++++++++--------- src/dynalab_core/dlpak.py | 67 ++++++++++++++++++++++++++++++++++++ src/dynalab_core/errors.py | 9 +++++ 4 files changed, 113 insertions(+), 17 deletions(-) create mode 100644 src/dynalab_core/dlpak.py diff --git a/src/dynalab_core/__init__.py b/src/dynalab_core/__init__.py index c97349d..3a08b95 100644 --- a/src/dynalab_core/__init__.py +++ b/src/dynalab_core/__init__.py @@ -2,6 +2,8 @@ # Copyright (C) 2026 Association Exergie # SPDX-License-Identifier: GPL-3.0-or-later + +from datetime import datetime, timezone import logging from queue import Empty, Queue import threading @@ -13,6 +15,7 @@ from uuid import UUID from dynalab_core.buffer import ValueBuffer from dynalab_core.config import CoreConfig from dynalab_core.constants import CORE_VERSION +from dynalab_core.dlpak import DLPak from dynalab_core.errors import CoreStateMismatchError from dynalab_core.protocols.endpoint import ConnectorRegistry from dynalab_core.protocols.common import VersionDescriptor @@ -34,6 +37,8 @@ class Core: ) self._recording = threading.Event() self._recording_buffer = ValueBuffer() + self._recording_timestamp: datetime = datetime.now(timezone.utc) + self._processing_buffer: DLPak | None = None self._core_version: VersionDescriptor = CORE_VERSION self._stop_event: threading.Event = threading.Event() self._core_config: CoreConfig = config @@ -106,11 +111,18 @@ class Core: def start_recording(self) -> None: self._recording_buffer.clear() + self._recording_timestamp = datetime.now(timezone.utc) self._recording.set() def stop_recording(self) -> None: self._recording.clear() - self._recording_buffer.export_csv("output.csv") + 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(), + ) + self._processing_buffer.write("./", "output") def get_signal_descriptor(self, signal_id: UUID) -> SignalDescriptor | None: return self._connector_registry.get_signal_descriptor(signal_id) diff --git a/src/dynalab_core/buffer.py b/src/dynalab_core/buffer.py index 6f5d511..0aec51f 100644 --- a/src/dynalab_core/buffer.py +++ b/src/dynalab_core/buffer.py @@ -4,6 +4,7 @@ from collections import defaultdict import csv +import io from pathlib import Path from threading import Lock from uuid import UUID @@ -22,11 +23,15 @@ class ValueBuffer: with self._lock: self._samples.clear() + def get_signal_ids(self) -> list[UUID]: + with self._lock: + return list({signal_id for _, signal_id, _ in self._samples}) + def __len__(self) -> int: with self._lock: return len(self._samples) - def export_csv(self, path: str | Path) -> None: + def export_csv(self) -> str: with self._lock: samples = list(self._samples) @@ -39,22 +44,25 @@ class ValueBuffer: ordered_signal_ids = sorted(signal_ids, key=str) - with Path(path).open(mode="w", newline="", encoding="utf-8") as file: - writer = csv.writer(file) + output = io.StringIO(newline="") + + writer = csv.writer(output) + + writer.writerow( + ["timestamp_ns", *(str(signal_id) for signal_id in ordered_signal_ids)] + ) + + for timestamp_ns in sorted(rows): + timestamp_values = rows[timestamp_ns] writer.writerow( - ["timestamp_ns", *(str(signal_id) for signal_id in ordered_signal_ids)] + [ + timestamp_ns, + *( + timestamp_values.get(signal_id, "") + for signal_id in ordered_signal_ids + ), + ] ) - for timestamp_ns in sorted(rows): - timestamp_values = rows[timestamp_ns] - - writer.writerow( - [ - timestamp_ns, - *( - timestamp_values.get(signal_id, "") - for signal_id in ordered_signal_ids - ), - ] - ) + return output.getvalue() diff --git a/src/dynalab_core/dlpak.py b/src/dynalab_core/dlpak.py new file mode 100644 index 0000000..be27118 --- /dev/null +++ b/src/dynalab_core/dlpak.py @@ -0,0 +1,67 @@ +# Copyright (C) 2026 Hector van der Aa +# Copyright (C) 2026 Association Exergie +# SPDX-License-Identifier: GPL-3.0-or-later + +from datetime import datetime, timezone +from pathlib import Path +from typing import Literal +from zipfile import ZIP_DEFLATED, ZipFile + +from pydantic import BaseModel, Field + +from dynalab_core.buffer import ValueBuffer +from dynalab_core.errors import DLPakNoDataError +from dynalab_core.protocols.packets.handshake import SignalDescriptor + + +class RecordManifest(BaseModel): + timestamp: datetime = Field(default_factory=lambda: datetime.now(timezone.utc)) + signals: list[SignalDescriptor] = Field(default_factory=list) + comments: str = "" + + +class DLPak: + def __init__(self) -> None: + self._dir: Path | None = None + self._filename: str | None = None + self._savetype: Literal["zip", "folder"] = "zip" + self._data: ValueBuffer | None = None + self._manifest: RecordManifest | None = None + + def set_data(self, data: ValueBuffer) -> None: + self._data = data + + def set_manifest( + self, timestamp: datetime, all_signals: list[SignalDescriptor] + ) -> None: + if self._data is None: + raise DLPakNoDataError + + data_signal_ids = self._data.get_signal_ids() + + manifest_signals: list[SignalDescriptor] = [] + + for signal in all_signals: + if signal.id in data_signal_ids: + manifest_signals.append(signal) + + self._manifest = RecordManifest( + timestamp=timestamp, signals=manifest_signals, comments="" + ) + + def write(self, dir: str | Path, filename: str) -> None: + output_dir = Path(dir) + + if output_dir.is_file(): + output_dir = output_dir.parent + + file_path = Path(filename) + + if file_path.suffix.lower() != ".dlpak": + file_path = file_path.with_suffix(".dlpak") + + output_path = output_dir / file_path + + with ZipFile(output_path, mode="w", compression=ZIP_DEFLATED) as archive: + archive.writestr("manifest.json", self._manifest.model_dump_json(indent=2)) + archive.writestr("data.csv", self._data.export_csv()) diff --git a/src/dynalab_core/errors.py b/src/dynalab_core/errors.py index 7e399c7..efa0c9e 100644 --- a/src/dynalab_core/errors.py +++ b/src/dynalab_core/errors.py @@ -9,3 +9,12 @@ class CoreError(Exception): class CoreStateMismatchError(CoreError): """DynaLab Core state error.""" + + +# DLPak error declarations +class DLPakError(Exception): + """DLPak error.""" + + +class DLPakNoDataError(DLPakError): + """DLPak no data error."""