From cdba587c5c33c5450944af25188a331f3d06337c Mon Sep 17 00:00:00 2001 From: Hector van der Aa Date: Mon, 6 Jul 2026 17:56:02 +0200 Subject: [PATCH] Initial dependency build --- src/dynalab_protocol/__init__.py | 5 ++- src/dynalab_protocol/models.py | 71 ++++++++++++++++++++++++++++++++ src/dynalab_protocol/wire.py | 33 +++++++++++++++ 3 files changed, 107 insertions(+), 2 deletions(-) create mode 100644 src/dynalab_protocol/models.py create mode 100644 src/dynalab_protocol/wire.py diff --git a/src/dynalab_protocol/__init__.py b/src/dynalab_protocol/__init__.py index d1b813a..1bd24e9 100644 --- a/src/dynalab_protocol/__init__.py +++ b/src/dynalab_protocol/__init__.py @@ -1,2 +1,3 @@ -def hello() -> str: - return "Hello from dynalab-protocol!" +# Copyright (C) 2026 Hector van der Aa +# Copyright (C) 2026 Association Exergie +# SPDX-License-Identifier: GPL-3.0-or-later diff --git a/src/dynalab_protocol/models.py b/src/dynalab_protocol/models.py new file mode 100644 index 0000000..e2dc3a7 --- /dev/null +++ b/src/dynalab_protocol/models.py @@ -0,0 +1,71 @@ +# Copyright (C) 2026 Hector van der Aa +# Copyright (C) 2026 Association Exergie +# SPDX-License-Identifier: GPL-3.0-or-later + +from typing import Annotated, Literal +from uuid import UUID +from pydantic import BaseModel, Field + +PROTOCOL_VERSION = 1 + + +class VersionDescriptor(BaseModel): + type: Literal["alpha", "beta", "release"] + major: int + minor: int + patch: int + + +class SignalDescriptor(BaseModel): + id: UUID + name: str + type: Literal["number", "binary"] + min_value: float | None = None + max_value: float | None = None + + +class DynaLabHello(BaseModel): + type: Literal["dynalab_hello"] = "dynalab_hello" + instance_id: UUID + version: VersionDescriptor + protocol_version: int = PROTOCOL_VERSION + + heartbeat_interval_ms: int = 1000 + heartbeat_timeout_ms: int = 5000 + + protocol_version: int = PROTOCOL_VERSION + + +class ConnectorHello(BaseModel): + type: Literal["connector_hello"] = "connector_hello" + connector_uuid: UUID + protocol_version: int = PROTOCOL_VERSION + + connector_name: str + connector_version: str + + signals: list[SignalDescriptor] + + +class HandshakeAccepted(BaseModel): + type: Literal["handshake_accepted"] = "handshake_accepted" + + accepted_signals: list[UUID] + + +class HandshakeRejected(BaseModel): + type: Literal["handshake_rejected"] = "handshake_rejected" + reason: str + + +class Heartbeat(BaseModel): + type: Literal["heartbeat"] = "heartbeat" + sequence: int + send_timestamp: int + return_timestamp: int | None = None + + +ProtocolMessage = Annotated[ + DynaLabHello | ConnectorHello | HandshakeAccepted | HandshakeRejected | Heartbeat, + Field(discriminator="type"), +] diff --git a/src/dynalab_protocol/wire.py b/src/dynalab_protocol/wire.py new file mode 100644 index 0000000..2bcc0e5 --- /dev/null +++ b/src/dynalab_protocol/wire.py @@ -0,0 +1,33 @@ +# Copyright (C) 2026 Hector van der Aa +# Copyright (C) 2026 Association Exergie +# SPDX-License-Identifier: GPL-3.0-or-later + +from asyncio import StreamReader, StreamWriter +from pydantic import TypeAdapter + +from dynalab_protocol.models import ProtocolMessage + + +_MESSAGE_ADAPTER = TypeAdapter(ProtocolMessage) + + +async def write_message(writer: StreamWriter, message: ProtocolMessage) -> None: + writer.write(message.model_dump_json().encode("utf-8") + b"\n") + await writer.drain() + + +async def read_message( + reader: StreamReader, *, max_frame_bytes: int = 65_536 +) -> ProtocolMessage: + line = await reader.readline() + + if not line: + raise ConnectionError("Peer disconnected") + + if len(line) > max_frame_bytes: + raise ValueError("Incoming protocol frame exceeds maximum size") + + if not line.endswith(b"\n"): + raise ValueError("Protocol frame is missing newline delimiter") + + return _MESSAGE_ADAPTER.validate_json(line)