diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 0cde6bb..1bef450 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -32,11 +32,12 @@ jobs: - run: python -m pytest python - if: matrix.python-version == '3.14' run: | - python -m ruff check python + python -m ruff check python benchmarks scripts python -m mypy python/src + python scripts/check_docs.py esp32: - name: ESP32 Arduino compile + name: ESP32 Arduino examples runs-on: ubuntu-latest steps: - uses: actions/checkout@v7 @@ -46,10 +47,25 @@ jobs: cache: pip - run: python -m pip install platformio - run: platformio run -e esp32dev + - run: platformio ci examples/SecureTelemetryClient/SecureTelemetryClient.ino --board esp32dev --lib src + - run: platformio ci examples/BluetoothSerialClient/BluetoothSerialClient.ino --board esp32dev --lib src + + arduino-native: + name: Arduino native tests + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v7 + - uses: actions/setup-python@v7 + with: + python-version: "3.14" + - run: python -m pip install platformio + - run: platformio test -e native package: name: Reproducible packages runs-on: ubuntu-latest + env: + SOURCE_DATE_EPOCH: "1700000000" steps: - uses: actions/checkout@v7 - uses: actions/setup-python@v7 @@ -57,12 +73,18 @@ jobs: python-version: "3.14" cache: pip - run: python -m pip install build - - run: python scripts/build_release.py - run: python -m build python + - run: python scripts/build_release.py + - run: python scripts/build_linux_bundle.py + - run: python scripts/build_checksums.py + - run: python scripts/verify_reproducible.py + - run: sh -n packaging/linux/install.sh packaging/linux/uninstall.sh - uses: actions/upload-artifact@v7 with: name: packages path: | dist/OpenNet-*.zip + dist/OpenNet-linux-*.tar.gz + dist/SHA256SUMS python/dist/* if-no-files-found: error diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 43e9184..f1f1d04 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -10,6 +10,8 @@ permissions: jobs: release: runs-on: ubuntu-latest + env: + SOURCE_DATE_EPOCH: "1700000000" steps: - uses: actions/checkout@v7 with: @@ -22,17 +24,24 @@ jobs: run: | python -m pip install -e "./python[dev]" build python -m pytest python - python -m ruff check python + python -m ruff check python benchmarks scripts python -m mypy python/src - python scripts/build_release.py --expected-version "${GITHUB_REF_NAME#v}" + python scripts/check_docs.py python -m build python + python scripts/build_release.py --expected-version "${GITHUB_REF_NAME#v}" + python scripts/build_linux_bundle.py --expected-version "${GITHUB_REF_NAME#v}" + python scripts/build_checksums.py + python scripts/verify_reproducible.py + sh -n packaging/linux/install.sh packaging/linux/uninstall.sh - name: Create GitHub release env: GH_TOKEN: ${{ github.token }} run: >- gh release create "$GITHUB_REF_NAME" dist/OpenNet-*.zip + dist/OpenNet-linux-*.tar.gz + dist/SHA256SUMS python/dist/* - --generate-notes + --notes-file "docs/releases/$GITHUB_REF_NAME.md" --verify-tag --title "OpenNet $GITHUB_REF_NAME" diff --git a/CHANGELOG.md b/CHANGELOG.md index d52aabb..6aadb8f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,31 @@ All notable changes are documented here. OpenNet follows ## [Unreleased] +## [0.1.1] - 2026-07-29 + +### Added + +- Unified `opennet` CLI with `serve`, `send`, `ping`, and `benchmark` commands. +- Acknowledgement retries using stable IDs and the ONP/1 `DUPLICATE` flag. +- Bounded server duplicate suppression, connection limits, operational counters, + receive queues, and connect timeouts. +- Mutual TLS options and secure-by-default remote CLI behavior. +- Generic Arduino `Stream` support for UART and Bluetooth Classic SPP transports. +- Typed Arduino receive accessors, readable error names, and ACK lookup helper. +- Hardened Linux/systemd installer bundle, checksums, and deterministic packages. +- Native C++ protocol tests, malformed-control tests, and retry/disconnect tests. +- Measured TCP/TLS/load benchmarks, protocol comparisons, tutorials, tips, + transport boundaries, security guidance, and honest adoption analytics. + +### Fixed + +- Partial Arduino transport writes now complete the whole frame. +- Python receivers now wake with an error when a peer disconnects. +- Handler errors no longer expose exception details to remote peers. +- Control-frame validation is consistent across Python and Arduino. + +## [0.1.0] - 2026-07-29 + ### Added - ONP/1 protocol specification and typed binary framing. @@ -13,4 +38,6 @@ All notable changes are documented here. OpenNet follows - Arduino IDE, PlatformIO, and Python examples. - Cross-platform tests, CI, packaging, and contributor documentation. -[Unreleased]: https://github.com/devkyato/OpenNet/compare/v0.1.0...HEAD +[Unreleased]: https://github.com/devkyato/OpenNet/compare/v0.1.1...HEAD +[0.1.1]: https://github.com/devkyato/OpenNet/compare/v0.1.0...v0.1.1 +[0.1.0]: https://github.com/devkyato/OpenNet/releases/tag/v0.1.0 diff --git a/README.md b/README.md index 07b4887..5668fdf 100644 --- a/README.md +++ b/README.md @@ -3,78 +3,122 @@ ![OpenNet cover](docs/assets/opennet-cover.png) [![CI](https://github.com/devkyato/OpenNet/actions/workflows/ci.yml/badge.svg)](https://github.com/devkyato/OpenNet/actions/workflows/ci.yml) +[![Release](https://img.shields.io/github/v/release/devkyato/OpenNet)](https://github.com/devkyato/OpenNet/releases) [![License: MIT](https://img.shields.io/badge/license-MIT-4a638f.svg)](LICENSE) [![Protocol](https://img.shields.io/badge/protocol-ONP%2F1-4a638f.svg)](docs/protocol.md) -OpenNet is a small, typed messaging protocol and library for sending JSON, text, -numbers, and arbitrary binary data between ESP32 devices, Raspberry Pi computers, -and backend services over Wi-Fi or Ethernet. +OpenNet is a dependency-free typed messaging protocol for direct communication +between ESP32 devices, Raspberry Pi computers, and backend services. It sends +JSON, text, signed integers, finite doubles, booleans, null, and arbitrary binary +data over TCP/TLS or another reliable ordered stream. + +Version 0.1.1 is an alpha release for real projects and learning. It has a +documented wire format, bounded resource use, retries and duplicate suppression, +native Arduino protocol tests, Python tests across 3.9–3.14, reproducible release +packages, measured local performance, and security guidance. It does not claim +that a radio or internet route exists when the underlying hardware/network is +unavailable. + +## Why OpenNet + +- One small API and one frame format across ESP32 and Python. +- Typed values without a JSON dependency for every scalar or binary payload. +- Optional application ACKs, retry with stable message IDs, and bounded + duplicate suppression. +- TCP/TLS over Wi-Fi, Ethernet, or routed networks. +- Generic Arduino `Stream` support for UART, USB serial, and Bluetooth Classic + SPP on compatible ESP32 hardware. +- Strict control frames, CRC-32 corruption detection, payload limits, bounded + queues, connection limits, and partial transport-write handling. +- CLI tools to serve, send, ping, and benchmark. +- Arduino IDE, PlatformIO/VS Code, Python, Linux/systemd, and security tutorials. + +## Install + +Python wheel from the release page: + +```sh +python -m pip install \ + https://github.com/devkyato/OpenNet/releases/download/v0.1.1/opennet_protocol-0.1.1-py3-none-any.whl +opennet --version +``` + +For an isolated command: -The first release deliberately targets one dependable path: an ESP32 client and a -Python 3.9+ client/server communicating over TCP. The wire format is documented and -language-neutral, so additional transports and language bindings can be added -without inventing incompatible message formats. +```sh +pipx install \ + https://github.com/devkyato/OpenNet/releases/download/v0.1.1/opennet_protocol-0.1.1-py3-none-any.whl +``` -## What it provides +For Raspberry Pi/Linux, extract `OpenNet-linux-0.1.1.tar.gz` and run +`sudo ./install.sh`. It creates a hardened, unprivileged systemd service with a +loopback-only default. For Arduino IDE, install `OpenNet-0.1.1.zip` through +**Sketch > Include Library > Add .ZIP Library**. -- A compact binary frame with versioning, message IDs, type information, CRC-32, - and a 16 MiB defensive payload limit. -- JSON, UTF-8 text, signed integers, IEEE-754 doubles, booleans, null, and raw bytes. -- Async Python client/server APIs for Raspberry Pi and backend systems. -- An Arduino-compatible ESP32 client with reconnect, heartbeat, and delivery ACKs. -- Optional TLS at the transport layer. -- Arduino IDE, PlatformIO, Python, and VS Code examples. -- Protocol conformance tests and cross-platform CI. +The Python distribution is installable with pip, but v0.1.1 is distributed from +GitHub Releases rather than the public PyPI index. -## Five-minute start +## Five-minute local test -Run the Python server: +Terminal one: -```bash -python -m pip install -e "./python[dev]" -opennet-server --host 0.0.0.0 --port 8765 +```sh +opennet serve --echo ``` -Install this repository as an Arduino library, open -`File > Examples > OpenNet > TelemetryClient`, set the Wi-Fi credentials and -server address, then upload it to an ESP32. +Terminal two: -Python clients are equally small: +```sh +opennet ping --count 5 +opennet send sensor/temperature 24.7 --type float +opennet send device/state '{"online":true}' --type json +opennet benchmark --count 1000 --payload-size 64 +``` + +Python code is equally small: ```python import asyncio from opennet import OpenNetClient async def main(): - async with OpenNetClient("192.168.1.50", 8765) as client: - await client.send("temperature", {"celsius": 24.7}) + async with OpenNetClient("127.0.0.1") as client: + message_id = await client.send( + "lab/temperature", 24.7, retries=2 + ) + print("acknowledged", message_id) asyncio.run(main()) ``` -See the [getting-started guide](docs/getting-started.md), [protocol -specification](docs/protocol.md), and [compatibility policy](docs/compatibility.md). - -## Project status - -OpenNet is pre-1.0 software. ONP/1 framing is stable for the v0.x series, while -higher-level APIs may improve based on real hardware feedback. Do not represent -untested boards or operating systems as supported; add a compatibility report when -you test one. - -## Contributing - -Student contributions are welcome. Good first tasks are labeled in the issue -tracker, and every feature should include tests or a reproducible hardware test -report. Read [CONTRIBUTING.md](CONTRIBUTING.md) and the -[code of conduct](CODE_OF_CONDUCT.md) before opening a pull request. - -## Security - -CRC detects accidental corruption; it is not encryption or authentication. Use TLS -on untrusted networks. See [SECURITY.md](SECURITY.md) for the threat model and -private reporting instructions. - -## License - -OpenNet is available under the [MIT License](LICENSE). +The CLI defaults to loopback and refuses remote plaintext unless it is explicitly +allowed for a trusted development network. Use verified TLS—and mutual TLS when +device identity matters—outside isolated local testing. + +## Evidence and design + +- [Measured TCP/TLS and 20-client load results](docs/benchmarks.md) +- [Comparison with raw TCP, MQTT, WebSocket, and HTTP](docs/comparison.md) +- [Architecture, ACK retry, and defensive boundaries](docs/architecture.md) +- [Actual GitHub adoption counters and analytics limits](docs/adoption-analytics.md) +- [Transport support, including Bluetooth boundaries](docs/transports.md) + +Measured loopback results are software baselines, not invented Wi-Fi or hardware +claims. Hardware field reports are welcome through the +[compatibility process](docs/compatibility.md). + +## Learn and build + +- [Tutorials](docs/tutorials.md) +- [CLI reference](docs/cli.md) +- [Linux/Raspberry Pi service guide](docs/linux-service.md) +- [Getting started](docs/getting-started.md) +- [Development tips](docs/tips.md) +- [Security guide](docs/security.md) +- [ONP/1 protocol specification](docs/protocol.md) +- [Contributing](CONTRIBUTING.md) + +Student contributions are welcome. Good changes include a reproducible test, +clear documentation, or a complete compatibility report. OpenNet follows +[Semantic Versioning](https://semver.org/) and is licensed under the +[MIT License](LICENSE). diff --git a/SECURITY.md b/SECURITY.md index 5f602d6..10c2050 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -25,3 +25,7 @@ coordinate disclosure after a fix is available. - Applications must authorize topics and validate payloads. - The 16 MiB frame limit reduces memory-exhaustion risk; deployments should choose a smaller application limit appropriate for their hardware. + +The CLI defaults to loopback and requires an explicit flag for remote plaintext. +Mutual TLS is available when deployments must authenticate client devices. See the +[security guide](docs/security.md) for commands, limits, and deployment checklists. diff --git a/benchmarks/loopback.py b/benchmarks/loopback.py new file mode 100644 index 0000000..0562770 --- /dev/null +++ b/benchmarks/loopback.py @@ -0,0 +1,136 @@ +"""Reproducible ONP/1 TCP/TLS loopback latency and load benchmark.""" + +from __future__ import annotations + +import argparse +import asyncio +import json +import platform +import ssl +import statistics +import sys +import time +from pathlib import Path +from typing import Any + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT / "python" / "src")) + +from opennet import OpenNetClient, OpenNetServer + + +def percentile(values: list[float], fraction: float) -> float: + ordered = sorted(values) + return ordered[round((len(ordered) - 1) * fraction)] + + +def summarize( + samples: list[float], payload_size: int, clients: int, elapsed: float +) -> dict[str, Any]: + return { + "payload_bytes": payload_size, + "clients": clients, + "messages": len(samples), + "mean_ms": round(statistics.fmean(samples), 3), + "p50_ms": round(percentile(samples, 0.50), 3), + "p95_ms": round(percentile(samples, 0.95), 3), + "max_ms": round(max(samples), 3), + "messages_per_second": round(len(samples) / elapsed, 1), + } + + +async def one_run( + *, + payload_size: int, + messages: int, + clients: int, + server_ssl: ssl.SSLContext | None, + client_ssl: ssl.SSLContext | None, +) -> dict[str, Any]: + server = OpenNetServer(lambda _peer, _frame: None, port=0, ssl_context=server_ssl) + await server.start() + address = server.sockets[0].getsockname() + port = int(address[1]) # type: ignore[index] + samples: list[float] = [] + payload = bytes(index % 251 for index in range(payload_size)) + + async def worker(count: int) -> None: + async with OpenNetClient( + "127.0.0.1", port, ssl_context=client_ssl, ack_timeout=10 + ) as client: + await client.send("_benchmark/warmup", payload) + for _ in range(count): + started = time.perf_counter() + await client.send("_benchmark/data", payload) + samples.append((time.perf_counter() - started) * 1000) + + counts = [messages // clients + (1 if index < messages % clients else 0) for index in range(clients)] + started = time.perf_counter() + await asyncio.gather(*(worker(count) for count in counts if count)) + elapsed = time.perf_counter() - started + await server.close() + result = summarize(samples, payload_size, clients, elapsed) + result["server_received"] = server.stats.received_messages + return result + + +def tls_contexts(args: argparse.Namespace) -> tuple[ssl.SSLContext | None, ssl.SSLContext | None]: + if args.cert is None: + return None, None + if args.key is None or args.ca is None: + raise SystemExit("--cert, --key, and --ca are required together") + server = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER) + server.load_cert_chain(args.cert, args.key) + client = ssl.create_default_context(cafile=str(args.ca)) + client.check_hostname = False + return server, client + + +async def run(args: argparse.Namespace) -> dict[str, Any]: + server_ssl, client_ssl = tls_contexts(args) + results = [] + for size in args.payload_sizes: + results.append( + await one_run( + payload_size=size, + messages=args.messages, + clients=args.clients, + server_ssl=server_ssl, + client_ssl=client_ssl, + ) + ) + return { + "schema": 1, + "transport": "TLS 1.2+" if server_ssl else "TCP", + "environment": { + "platform": platform.platform(), + "python": platform.python_version(), + "processor": platform.processor() or "not reported", + }, + "method": "IPv4 loopback; one warm-up message per client; ACK round-trip timing", + "results": results, + } + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--messages", type=int, default=1000) + parser.add_argument("--clients", type=int, default=1) + parser.add_argument("--payload-sizes", type=int, nargs="+", default=[0, 64, 1024, 65536]) + parser.add_argument("--cert", type=Path) + parser.add_argument("--key", type=Path) + parser.add_argument("--ca", type=Path) + parser.add_argument("--output", type=Path) + args = parser.parse_args() + if args.messages <= 0 or args.clients <= 0 or any(size < 0 for size in args.payload_sizes): + parser.error("messages and clients must be positive; payload sizes cannot be negative") + text = json.dumps(asyncio.run(run(args)), indent=2, sort_keys=True) + "\n" + if args.output: + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(text, encoding="utf-8") + else: + print(text, end="") + + +if __name__ == "__main__": + main() diff --git a/docs/adoption-analytics.md b/docs/adoption-analytics.md new file mode 100644 index 0000000..c0075e3 --- /dev/null +++ b/docs/adoption-analytics.md @@ -0,0 +1,32 @@ +# Adoption analytics + +OpenNet does not collect telemetry from library users. The only current adoption +signals are aggregate GitHub repository and release counters. + +Snapshot queried from the GitHub API on 2026-07-29: + +| Signal | Actual value | +|---|---:| +| GitHub views in returned 14-day window | 0 | +| Unique viewers in returned window | 0 | +| Git clones in returned window | 1 | +| Unique cloners in returned window | 1 | +| Stars | 0 | +| Forks | 0 | +| Watchers/subscribers | 1 | +| Open issues | 3 | +| v0.1.0 release asset downloads | 0 | + +```mermaid +xychart-beta + title "Current public adoption signals (2026-07-29 snapshot)" + x-axis ["Unique clones", "Subscribers", "Stars", "Forks", "Downloads"] + y-axis "Count" 0 --> 1 + bar [1, 1, 0, 0, 0] +``` + +These values are a transparent early-project baseline, not evidence of a user +population. GitHub traffic is a rolling window and may lag. Demographic data +(age, location, school, occupation, or device ownership) is unavailable and +should not be guessed. Future releases can report platform compatibility and +opt-in survey results without adding tracking to the library. diff --git a/docs/architecture.md b/docs/architecture.md index e121881..13114a1 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -1,44 +1,55 @@ # Architecture -OpenNet separates the parts that must interoperate from platform-specific I/O. - -```text -Application values and topics - | - OpenNet API - | - ONP/1 frame codec <- language-neutral and tested with vectors - | - TCP or TLS byte stream <- reliable ordered transport - | - Wi-Fi / Ethernet / OS networking +OpenNet separates ONP/1 framing from the byte stream carrying it. Both endpoints +must provide a reliable, ordered stream; Wi-Fi itself is not required. + +```mermaid +flowchart LR + A["Application values
JSON · text · numbers · bytes"] --> B["OpenNet API"] + B --> C["ONP/1 codec
24-byte header · type · topic · CRC-32"] + C --> D{"Reliable ordered stream"} + D --> E["TCP/TLS
Wi-Fi, Ethernet, internet"] + D --> F["Arduino Stream
UART, USB serial"] + D --> G["Bluetooth Classic SPP
supported ESP32 models"] ``` -The Python package provides the reference codec plus asyncio client and server. The -Arduino library uses the ESP32 networking clients and a bounded receive state -machine. Both implementations use the same 24-byte header. - -## Reliability - -TCP already retransmits lost packets and preserves order. ONP/1 acknowledgements -add application-level evidence that a complete frame was parsed and accepted. -Applications needing durable delivery must persist outgoing messages and decide -how to handle duplicates after reconnecting; OpenNet does not pretend RAM queues -survive power loss. - -## Latency +The Python package provides the reference codec and asyncio TCP/TLS client and +server. The Arduino implementation accepts either a `Client` (`WiFiClient`, +`WiFiClientSecure`) or an already-open Arduino `Stream` (`Serial`, +`BluetoothSerial`). The frame bytes are identical across transports. + +BLE GATT is packet-oriented, not a continuous byte stream. It needs an adapter +that fragments and reassembles ONP frames; v0.1.1 does not claim direct BLE +support. ESP32-C3/S3 boards also do not provide Bluetooth Classic SPP. + +## Delivery and retry + +```mermaid +sequenceDiagram + participant C as Client + participant S as Server + C->>S: DATA id=42, ACK_REQUIRED + Note over C: ACK timeout + C->>S: DATA id=42, ACK_REQUIRED + DUPLICATE + Note over S: bounded recent-ID window + S-->>C: ACK id=42, DUPLICATE + Note over C: delivery confirmed +``` -OpenNet disables Nagle's algorithm where supported, uses a compact header, and -does not require a broker round trip. Actual latency still depends on radio -conditions, access-point load, power-saving mode, operating-system scheduling, and -the application. Low-signal or disconnected areas cannot be made reliable by a -software library alone. +Python retries reuse the message ID, and the server suppresses duplicate handler +execution within a configurable bounded window. An ACK proves that the peer +parsed and accepted the frame; it does not make delivery durable across power +loss. Durable applications must persist outbound messages and application state. -## Planned extensions +## Defensive boundaries -- Authenticated session handshake and topic authorization helpers. -- Optional service discovery for trusted LANs. -- Persistent outbound queues. -- TypeScript and portable C bindings. +- Payload, topic, receive-queue, connection, and duplicate-window bounds prevent + unbounded routine growth. +- The CRC-32 catches accidental corruption before the application sees a value. +- Strict type and control-frame validation rejects malformed inputs. +- Remote CLI connections require TLS unless plaintext is explicitly enabled. +- TLS provides confidentiality and server identity; mutual TLS can also identify + clients. -Extensions will be proposed in issues and specified before implementation. +See [Security](security.md) for the trust model and [Transports](transports.md) +for supported connection types. diff --git a/docs/benchmarks.md b/docs/benchmarks.md new file mode 100644 index 0000000..8c5d10b --- /dev/null +++ b/docs/benchmarks.md @@ -0,0 +1,45 @@ +# Measured performance + +These are actual controlled measurements from the v0.1.1 worktree on +2026-07-29. They are reproducible with `benchmarks/loopback.py`; they are not +estimates and are not presented as Wi-Fi, Bluetooth, Raspberry Pi, or ESP32 +field results. + +Environment: Windows 11 build 26200, AMD64 Family 25 Model 68, CPython 3.14.4, +IPv4 loopback. Each row uses one warm-up message per client and then times a +complete DATA-to-ACK round trip. TLS used a locally generated RSA-2048 +certificate with certificate validation enabled. + +| Transport | Clients | Payload | Messages | p50 | p95 | Mean | Throughput | +|---|---:|---:|---:|---:|---:|---:|---:| +| TCP | 1 | 0 B | 1,000 | 0.164 ms | 0.279 ms | 0.195 ms | 5,069.6 msg/s | +| TCP | 1 | 64 B | 1,000 | 0.156 ms | 0.264 ms | 0.178 ms | 5,566.1 msg/s | +| TCP | 1 | 1 KiB | 1,000 | 0.159 ms | 0.265 ms | 0.186 ms | 5,346.0 msg/s | +| TCP | 1 | 64 KiB | 1,000 | 0.224 ms | 0.329 ms | 0.245 ms | 4,045.8 msg/s | +| TLS | 1 | 0 B | 1,000 | 0.231 ms | 0.351 ms | 0.262 ms | 3,578.0 msg/s | +| TLS | 1 | 64 B | 1,000 | 0.226 ms | 0.340 ms | 0.254 ms | 3,843.5 msg/s | +| TLS | 1 | 1 KiB | 1,000 | 0.236 ms | 0.357 ms | 0.265 ms | 3,712.9 msg/s | +| TLS | 1 | 64 KiB | 1,000 | 0.361 ms | 0.572 ms | 0.434 ms | 2,270.7 msg/s | +| TCP load | 20 | 64 B | 2,000 | 2.879 ms | 5.296 ms | 3.177 ms | 5,861.5 msg/s | + +```mermaid +xychart-beta + title "Single-client p95 ACK latency on this host" + x-axis ["0 B", "64 B", "1 KiB", "64 KiB"] + y-axis "Milliseconds" 0 --> 0.6 + line "TCP" [0.279, 0.264, 0.265, 0.329] + line "TLS" [0.351, 0.340, 0.357, 0.572] +``` + +## Interpretation and limits + +The run proves that the codec, client, server, ACK path, concurrency handling, +and verified TLS path work together on this local device under the stated load. +Loopback removes radio interference and network hops, so these numbers are a +software baseline—not a prediction for a real deployment. Wi-Fi latency depends +on signal, access-point congestion, power saving, and hardware. Bluetooth SPP +depends on radio and serial baud/stack behavior. + +For a field report, record the exact board, firmware, access point, distance/RSSI, +payload mix, TLS mode, run duration, reconnects, failures, and p50/p95/p99. The +[compatibility guide](compatibility.md) contains a report checklist. diff --git a/docs/cli.md b/docs/cli.md new file mode 100644 index 0000000..cc88934 --- /dev/null +++ b/docs/cli.md @@ -0,0 +1,54 @@ +# Command-line reference + +`opennet --help` exposes four commands. + +## Serve + +```sh +opennet serve [--host HOST] [--port PORT] [--echo] +``` + +Useful controls include `--max-payload`, `--max-connections`, +`--deduplication-window`, `--tls-cert`, `--tls-key`, and `--tls-client-ca`. +Remote plaintext listeners require `--allow-plaintext`. + +## Send + +```sh +opennet send TOPIC [VALUE] [--type TYPE] +``` + +Types are `text`, `json`, `int`, `float`, `bool`, `null`, `hex`, and `file`. +The command waits for an ACK by default and retries twice. Adjust with +`--retries`, `--retry-delay`, `--ack-timeout`, or `--no-ack`. + +Examples: + +```sh +opennet send status online +opennet send telemetry '{"temperature":24.7}' --type json +opennet send firmware/chunk ./part.bin --type file +opennet send counter 42 --type int --host pi.local --ca lab-ca.crt +``` + +## Ping and benchmark + +```sh +opennet ping --count 10 +opennet benchmark --count 1000 --warmup 25 --payload-size 1024 --json +``` + +Both commands measure complete ONP round trips rather than ICMP. `benchmark` +uses acknowledged binary DATA frames. + +## Connection and TLS options + +- `--host` defaults to `OPENNET_HOST` or `127.0.0.1`. +- `--port` defaults to `OPENNET_PORT` or `8765`. +- `--ca` enables TLS and verifies the server against that CA. +- `--cert` and `--key` add a client identity for mutual TLS. +- `--insecure` encrypts without identity verification and is diagnostic only. +- `--allow-plaintext` deliberately permits unencrypted remote traffic. + +The service also reads `OPENNET_MAX_PAYLOAD`. Command-line values take priority +over environment defaults. diff --git a/docs/comparison.md b/docs/comparison.md new file mode 100644 index 0000000..53ee2aa --- /dev/null +++ b/docs/comparison.md @@ -0,0 +1,75 @@ +# Protocol comparison + +OpenNet is a direct typed messaging option, not a universal replacement for +established web or broker protocols. + +| Capability | OpenNet ONP/1 | Raw TCP | MQTT 5 | WebSocket | HTTP | +|---|---|---|---|---|---| +| Primary topology | Direct client/server | Whatever the app invents | Brokered pub/sub | Direct full-duplex | Request/response | +| Message boundary | Built in | None | Built in | Built in | Built in | +| App value types | 7 typed values | None | Application-defined bytes | Text/binary frames | Media type/body | +| Topics | Built in | None | Built in | Application-defined | URI/application-defined | +| Delivery signal | Optional ONP ACK | None at app layer | QoS 0/1/2 | None at app layer | Response status | +| Browser-native | No | No | Usually via WebSocket | Yes | Yes | +| Infrastructure | One peer listens | Application-defined | Broker | WebSocket server | HTTP server | +| Best fit | Small direct MCU/Python messaging | Custom low-level protocol | Fleet pub/sub and offline sessions | Browser full-duplex apps | Web APIs and integrations | + +## Framing overhead + +```mermaid +xychart-beta + title "Minimum application framing bytes (topic/body excluded)" + x-axis ["Raw TCP", "WebSocket", "OpenNet"] + y-axis "Bytes" 0 --> 24 + bar [0, 2, 24] +``` + +This graph is structural, not a speed benchmark. Raw TCP adds no application +framing and therefore supplies no message boundary. A WebSocket frame starts at +2 bytes and grows with payload length; client-to-server frames also carry a +4-byte masking key. OpenNet uses a fixed 24-byte header and then its UTF-8 topic +and payload. MQTT and HTTP overhead vary too much by properties and headers to +represent as one honest number. + +## Choosing deliberately + +- Choose **MQTT** when you need a broker, retained messages, subscriptions, or + standardized QoS across a device fleet. +- Choose **WebSocket** when browser support is central. +- Choose **HTTP** when existing web tooling, caching, proxies, and REST semantics + matter more than a compact persistent channel. +- Choose **raw TCP** when you are prepared to design, test, and maintain every + application framing and typing rule. +- Choose **OpenNet** for direct, dependency-light typed messages shared by ESP32 + and Python code, with a small API and documented frame format. + +Primary specifications: [MQTT 5.0 (OASIS)](https://docs.oasis-open.org/mqtt/mqtt/v5.0/mqtt-v5.0.html), +[WebSocket RFC 6455](https://www.rfc-editor.org/rfc/rfc6455.html), and +[HTTP Semantics RFC 9110](https://www.rfc-editor.org/rfc/rfc9110.html). + +## Ecosystem usage indicator + +This is an actual dated package-download indicator, not protocol market share. +PyPI Stats reported the following last-month downloads on 2026-07-29: + +| Representative Python package | Protocol/use | Last-month downloads | +|---|---|---:| +| `paho-mqtt` | MQTT client | 8,134,970 | +| `websockets` | WebSocket library | 503,575,188 | +| `requests` | HTTP client | 1,678,041,738 | +| `opennet-protocol` | OpenNet | Not on PyPI; GitHub v0.1.0 assets had 0 downloads | + +```mermaid +xychart-beta + title "PyPI ecosystem indicator: log10(last-month downloads + 1)" + x-axis ["OpenNet", "paho-mqtt", "websockets", "requests"] + y-axis "log10 downloads" 0 --> 10 + bar [0, 6.91, 8.70, 9.22] +``` + +Established ecosystems are many orders of magnitude larger, which means more +integrations, operational knowledge, and independent testing. The counts also +include automation and repeat installs and exclude implementations in other +languages, so they must not be interpreted as unique people or deployments. +Source and method: [PyPI Stats API](https://pypistats.org/api/), which aggregates +public PyPI download data and excludes known mirrors. diff --git a/docs/compatibility.md b/docs/compatibility.md index dde9524..f7e5c4b 100644 --- a/docs/compatibility.md +++ b/docs/compatibility.md @@ -4,28 +4,36 @@ | Platform | Minimum | Status | Verification | |---|---|---|---| -| Python | CPython 3.9 | Supported | CI on 3.9–3.13 | -| Raspberry Pi OS | Python 3.9+ | Supported | Same pure-Python code; hardware reports invited | -| ESP32 Arduino core | 2.x, 3.x | Supported | CI compile for `esp32dev`; hardware reports invited | -| VS Code | Current + PlatformIO | Supported workflow | PlatformIO configuration | -| Arduino IDE | 2.x | Supported workflow | `library.properties` and examples | +| Python | CPython 3.9 | Supported | CI on 3.9–3.14 | +| Raspberry Pi OS | Python 3.9+ | Supported | Pure Python and Linux bundle; hardware reports invited | +| ESP32 Arduino | Current PlatformIO package | Supported alpha | `esp32dev` compile and native protocol tests | +| VS Code | Current + PlatformIO | Supported workflow | Checked-in PlatformIO configuration | +| Arduino IDE | 2.x | Supported workflow | Library metadata and examples | + +“Raspberry Pi support” means any model and operating-system image able to run a +supported CPython version and provide the chosen connection. It does not imply +that every historical Pi image, radio, or adapter has been physically tested. ## Not claimed -- Original ESP8266, RP2040/Pico W, Arduino UNO WiFi, or non-ESP32 Arduino boards. -- Every historical Raspberry Pi OS image. A board must run a maintained Python - 3.9+ environment and have working TCP/IP networking. -- Guaranteed operation through NAT, captive portals, or without network coverage. -- Real-time deadlines or safety certification. +- ESP8266, RP2040/Pico W, Arduino UNO WiFi, or non-ESP32 Arduino boards. +- Every historical Raspberry Pi OS image. +- Direct BLE GATT support. Bluetooth Classic SPP works only on ESP32 targets that + provide it; ESP32-C3/S3 do not. +- Guaranteed operation through NAT, captive portals, firewalls, or without + network coverage. +- Real-time deadlines, independent security audit, or safety certification. -## Adding a report +## Adding a field report Open a compatibility issue with: -- exact device and revision; -- OS, Python, Arduino core, and toolchain versions; -- Wi-Fi or Ethernet adapter; -- client/server example used; -- duration, message count, payload sizes, and observed failures. +- exact device, revision, OS, Python/Arduino core, and toolchain; +- Wi-Fi/Ethernet/Bluetooth adapter and access point; +- connection type and TLS mode; +- example or application commit used; +- duration, message count, payload sizes, and reconnect/failure count; +- RSSI/range and p50/p95/p99 latency when applicable. -Verified reports may be promoted into the supported table. +Verified reports may be promoted into the supported table. A failed test is useful +evidence too; do not omit it. diff --git a/docs/getting-started.md b/docs/getting-started.md index eedf3e5..d5dbb78 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -13,19 +13,15 @@ This includes current Raspberry Pi OS releases on Pi boards capable of running a supported Python version. ```bash -git clone https://github.com/devkyato/OpenNet.git -cd OpenNet -python -m venv .venv -source .venv/bin/activate -python -m pip install -e "./python" -opennet-server --host 0.0.0.0 --port 8765 +python -m pip install \ + https://github.com/devkyato/OpenNet/releases/download/v0.1.1/opennet_protocol-0.1.1-py3-none-any.whl +opennet serve ``` On Windows, activate with `.venv\Scripts\activate`. -The example server logs metadata, not full binary payloads. For internet-facing -deployments, configure TLS and an application authorization policy before opening -the firewall. +The default is loopback-only. Follow the [security guide](security.md) to configure +verified or mutual TLS before exposing the service. ## ESP32 with Arduino IDE @@ -52,7 +48,7 @@ For another project: ```ini lib_deps = - https://github.com/devkyato/OpenNet.git#v0.1.0 + https://github.com/devkyato/OpenNet.git#v0.1.1 ``` ## Sending values @@ -70,3 +66,6 @@ client.sendBool("door/open", false); Validate incoming JSON against an application schema and bound binary payload sizes for the receiving device. + +Continue with the [step-by-step tutorials](tutorials.md) and +[development tips](tips.md). diff --git a/docs/linux-service.md b/docs/linux-service.md new file mode 100644 index 0000000..19b2b18 --- /dev/null +++ b/docs/linux-service.md @@ -0,0 +1,70 @@ +# Linux and Raspberry Pi service + +The release bundle supports systemd systems with Python 3.9+ and the standard +library `venv` module. + +## Install + +```sh +grep 'OpenNet-linux-0.1.1.tar.gz' SHA256SUMS | sha256sum -c - +tar -xzf OpenNet-linux-0.1.1.tar.gz +cd OpenNet-linux-0.1.1 +sudo ./install.sh +``` + +The installer creates: + +- `/opt/opennet/venv` with the bundled wheel; +- an unprivileged `opennet` system account; +- `/etc/opennet/opennet.env`; +- `/etc/systemd/system/opennet.service`. + +The default listener is `127.0.0.1:8765`. It can be used immediately by local +applications and is not reachable from another machine. + +## Configure verified TLS for a LAN + +Install certificates readable by the service: + +```sh +sudo install -d -m 0750 -o root -g opennet /etc/opennet/tls +sudo install -m 0640 -o root -g opennet server.crt server.key \ + /etc/opennet/tls/ +``` + +Create a systemd override: + +```sh +sudo systemctl edit opennet +``` + +```ini +[Service] +ExecStart= +ExecStart=/opt/opennet/venv/bin/opennet serve --host 0.0.0.0 --tls-cert /etc/opennet/tls/server.crt --tls-key /etc/opennet/tls/server.key +``` + +Then: + +```sh +sudo systemctl daemon-reload +sudo systemctl restart opennet +sudo systemctl status opennet +``` + +For mutual TLS, install the device CA and append +`--tls-client-ca /etc/opennet/tls/devices-ca.crt`. + +## Operate and upgrade + +```sh +journalctl -u opennet -f +systemctl show opennet -p ActiveState -p SubState -p NRestarts +``` + +To upgrade, extract the new bundle and rerun `sudo ./install.sh`; the installer +preserves `/etc/opennet/opennet.env`. Review release notes and restart the service. + +To remove code while keeping configuration, run `sudo ./uninstall.sh`. Use +`sudo ./uninstall.sh --purge` only when the configuration and service account +should also be removed. diff --git a/docs/release.md b/docs/release.md index db8a89c..c717ed7 100644 --- a/docs/release.md +++ b/docs/release.md @@ -5,10 +5,21 @@ `python/pyproject.toml`. 3. Create and push an annotated `vX.Y.Z` tag. 4. The release workflow verifies that all versions match, runs tests, builds the - Python wheel/source distribution and Arduino ZIP, then creates a GitHub release. + Python wheel/source distribution, Arduino ZIP, Linux/systemd bundle, and + SHA-256 checksums, then creates a GitHub release. 5. A maintainer checks installation from the published artifacts and edits the generated release notes if needed. PyPI publication is intentionally not automatic until the owner configures trusted publishing. Release assets remain installable and reproducible without a registry credential. + +## v0.1.1 assets + +| Asset | Use | +|---|---| +| `OpenNet-0.1.1.zip` | Arduino IDE library | +| `opennet_protocol-0.1.1-py3-none-any.whl` | Direct pip/pipx install | +| `opennet_protocol-0.1.1.tar.gz` | Python source distribution | +| `OpenNet-linux-0.1.1.tar.gz` | `sudo` Linux/systemd install | +| `SHA256SUMS` | Artifact integrity verification | diff --git a/docs/releases/v0.1.1.md b/docs/releases/v0.1.1.md new file mode 100644 index 0000000..b6371a0 --- /dev/null +++ b/docs/releases/v0.1.1.md @@ -0,0 +1,52 @@ +# OpenNet v0.1.1 + +OpenNet v0.1.1 turns the initial ONP/1 implementation into a usable alpha for +direct ESP32, Raspberry Pi, and backend messaging. + +## Highlights + +- Reliable Python delivery retries preserve message IDs and mark retransmissions, + while servers suppress recent duplicates with bounded memory. +- The `opennet` CLI can serve, send every supported type, ping, and benchmark. +- Remote CLI traffic requires TLS unless plaintext is deliberately enabled. + Verified TLS and mutual TLS are both supported. +- Arduino now completes partial writes, validates control frames, decodes typed + values, and works with generic `Stream` transports such as UART and Bluetooth + Classic SPP on compatible ESP32 boards. +- A hardened Linux/systemd bundle provides a simple `sudo ./install.sh` path. +- Tutorials, tips, architecture diagrams, protocol comparisons, adoption + counters, and measured TCP/TLS/load results are included. + +## Install + +```sh +python -m pip install \ + https://github.com/devkyato/OpenNet/releases/download/v0.1.1/opennet_protocol-0.1.1-py3-none-any.whl +``` + +Arduino IDE users should install `OpenNet-0.1.1.zip`. Linux/Raspberry Pi users can +extract `OpenNet-linux-0.1.1.tar.gz` and run `sudo ./install.sh`. Verify downloads +against `SHA256SUMS`. + +## Verification + +- 41 Python tests locally on CPython 3.14.4; CI covers CPython 3.9–3.14. +- Native C++ protocol suite passed. +- Telemetry, verified-TLS, and Bluetooth Classic SPP examples compiled for an + ESP32 Dev Module. +- Fresh-wheel install, CLI version, live server ping, and typed-send smoke tests. +- Reproducibility checks for Arduino and Linux bundles. + +The measured loopback baseline and exact limitations are documented in +[`docs/benchmarks.md`](https://github.com/devkyato/OpenNet/blob/v0.1.1/docs/benchmarks.md). +It is not presented as ESP32 or Wi-Fi field performance. + +## Security and compatibility + +CRC-32 is not encryption or authentication. Use verified TLS outside isolated +local tests, use mutual TLS when device identity matters, and authorize topics in +the application. BLE GATT is not directly supported in this release; Bluetooth +Classic SPP requires compatible ESP32 hardware. + +See the full [changelog](https://github.com/devkyato/OpenNet/blob/v0.1.1/CHANGELOG.md) +and [security guide](https://github.com/devkyato/OpenNet/blob/v0.1.1/docs/security.md). diff --git a/docs/security.md b/docs/security.md new file mode 100644 index 0000000..3dd43c2 --- /dev/null +++ b/docs/security.md @@ -0,0 +1,56 @@ +# Security guide + +OpenNet is secure only when its deployment matches its threat model. CRC-32 is an +integrity check for accidental corruption; it does not stop an attacker. + +## Recommended modes + +| Environment | Minimum mode | +|---|---| +| One-process tests | Plain loopback | +| Trusted, isolated classroom LAN | TLS preferred; plaintext only by explicit choice | +| Shared Wi-Fi, campus, or internet | Verified TLS | +| Devices that must be individually trusted | Mutual TLS plus topic authorization | +| Safety-critical control | Not supported without an independent safety design | + +The CLI refuses non-loopback plaintext unless `--allow-plaintext` is supplied. +The server defaults to `127.0.0.1`. `--insecure` still encrypts but disables +certificate verification and is for diagnosis only. + +## Verified TLS + +Server: + +```sh +opennet serve --host 0.0.0.0 \ + --tls-cert server.crt --tls-key server.key +``` + +Client: + +```sh +opennet send sensor/temperature 24.7 --type float \ + --host opennet.example.net --ca lab-ca.crt +``` + +For mutual TLS, add `--tls-client-ca devices-ca.crt` to the server and +`--cert device.crt --key device.key` to the client. Issue a distinct certificate +per device so one device can be revoked without replacing every key. + +## Application responsibilities + +- Authorize which device may publish or receive each topic. +- Validate decoded values, JSON schemas, ranges, and file formats. +- Keep private keys outside the repository and restrict filesystem permissions. +- Rotate certificates and define a revocation process. +- Persist operation IDs when duplicate physical actions would be dangerous. +- Rate-limit expensive handlers and monitor `ServerStats`. +- Keep payload and connection limits appropriate for available RAM and CPU. + +## Current limits + +ONP/1 has no message-level signature or session authorization handshake. +TLS protects network streams, but a peer with an accepted certificate can still +send any topic unless the application checks it. Bluetooth pairing is not a +replacement for application authorization. OpenNet has not been independently +audited or certified. diff --git a/docs/tips.md b/docs/tips.md new file mode 100644 index 0000000..83d343a --- /dev/null +++ b/docs/tips.md @@ -0,0 +1,21 @@ +# Development tips + +- Start on loopback (`127.0.0.1`) and run `opennet ping` before adding a radio. +- Use short, stable topics such as `lab1/sensor3/temperature`; put changing data + in the value, not the topic. +- Set the smallest realistic `max_payload`. An ESP32 should not inherit a backend + server's multi-megabyte allowance. +- Call Arduino `poll()` every loop iteration. Long `delay()` calls postpone ACKs, + pings, and incoming messages. +- Retry idempotent operations. For actions such as “unlock door,” include an + application operation ID and persist processed IDs. +- Use JSON for evolving records, typed numbers for frequent scalar telemetry, and + bytes for already encoded data. +- Do not base64-encode binary data; ONP carries bytes directly. +- Use `opennet benchmark` on the deployment network and report p95/p99, not only + the fastest observation. +- Treat an ACK as peer acceptance, not durable storage or physical action. +- Validate JSON shape, topic authorization, ranges, and binary format in the + application handler. +- Never publish private keys, Wi-Fi passwords, addresses, or production + certificates in an issue or example. diff --git a/docs/transports.md b/docs/transports.md new file mode 100644 index 0000000..a3fc6a0 --- /dev/null +++ b/docs/transports.md @@ -0,0 +1,41 @@ +# Transports + +ONP/1 needs a reliable, ordered sequence of bytes. It does not require a specific +radio. + +| Connection | v0.1.1 path | Security | Notes | +|---|---|---|---| +| Wi-Fi/Ethernet/LAN | TCP | TLS recommended | Python and ESP32 | +| Internet/VPN | TCP | Verified TLS required by CLI | Routers/firewalls still apply | +| ESP32 Wi-Fi | `WiFiClient` / `WiFiClientSecure` | Prefer `WiFiClientSecure` | All targets need a compatible ESP32 Arduino core | +| UART/USB serial | Arduino `Stream` | Physical/access controls | Already-open stream | +| Bluetooth Classic SPP | ESP32 `BluetoothSerial` stream | Pairing plus app controls | Original ESP32; not C3/S3 | +| BLE GATT | Not direct in v0.1.1 | Pairing is not app authorization | Needs fragmentation/reassembly adapter | + +An internet connection does not make arbitrary devices discoverable or reachable. +Peers still need an address, route, open firewall, and listening service. Offline +radio links work only while both devices are in range. + +## Arduino stream use + +```cpp +#include +#include + +BluetoothSerial radio; +OpenNetClient openNet(radio); + +void setup() { + radio.begin("OpenNet-ESP32"); +} + +void loop() { + openNet.poll(); + // Once the SPP peer is connected: + // openNet.sendText("device/status", "online"); +} +``` + +The same constructor works with `Serial2` after the application calls +`Serial2.begin(...)`. For network clients, keep using `OpenNetClient(WiFiClient&)` +and `connect(host, port)`. diff --git a/docs/tutorials.md b/docs/tutorials.md new file mode 100644 index 0000000..1b45e88 --- /dev/null +++ b/docs/tutorials.md @@ -0,0 +1,73 @@ +# Tutorials + +## 1. Local Python round trip + +Install the v0.1.1 wheel directly from the GitHub release: + +```sh +python -m venv .venv +source .venv/bin/activate +python -m pip install \ + https://github.com/devkyato/OpenNet/releases/download/v0.1.1/opennet_protocol-0.1.1-py3-none-any.whl +``` + +Terminal one: + +```sh +opennet serve --echo +``` + +Terminal two: + +```sh +opennet ping --count 5 +opennet send sensor/temperature 24.7 --type float +opennet send device/state '{"online":true,"battery":91}' --type json +``` + +## 2. Raspberry Pi service + +Download and extract `OpenNet-linux-0.1.1.tar.gz`, then: + +```sh +sudo ./install.sh +systemctl status opennet +journalctl -u opennet -f +``` + +It listens only on `127.0.0.1` by default. Use `--no-start` if you need to add +TLS before the first start. See [Security](security.md). + +## 3. ESP32 over a trusted development LAN + +1. Install `OpenNet-0.1.1.zip` through Arduino IDE's **Add .ZIP Library**. +2. Open **File > Examples > OpenNet > TelemetryClient**. +3. Enter the Wi-Fi details and the server's LAN address. +4. Start a development server with + `opennet serve --host 0.0.0.0 --allow-plaintext`. +5. Upload and watch the serial monitor. + +This plaintext mode is for a controlled lab. Use the +`SecureTelemetryClient` example and a CA certificate for shared networks. + +## 4. ESP32 Bluetooth Classic SPP + +Open `BluetoothSerialClient` on two original ESP32 boards (or pair one to a host +serial port). ONP framing runs over `BluetoothSerial`, so JSON, numbers, bytes, +ACKs, and CRC validation remain identical. + +ESP32-C3 and ESP32-S3 do not support Bluetooth Classic SPP. BLE GATT requires a +future adapter and must not be described as supported by this release. + +## 5. Measure your real network + +Start the server and run: + +```sh +opennet benchmark --host raspberrypi.local --ca lab-ca.crt \ + --count 1000 --warmup 25 --payload-size 64 --json +``` + +Repeat at different distances and busy periods. Record device versions, RSSI, +payload, sample count, reconnects, p50, p95, and maximum. A good report includes +failures and configuration, not only the best number. diff --git a/examples/BluetoothSerialClient/BluetoothSerialClient.ino b/examples/BluetoothSerialClient/BluetoothSerialClient.ino new file mode 100644 index 0000000..d7d452a --- /dev/null +++ b/examples/BluetoothSerialClient/BluetoothSerialClient.ino @@ -0,0 +1,30 @@ +#include +#include + +#if !defined(CONFIG_BT_SPP_ENABLED) +#error Bluetooth Classic SPP is unavailable on this ESP32 target +#endif + +BluetoothSerial radio; +OpenNetClient openNet(radio); +unsigned long lastSend = 0; + +void setup() { + Serial.begin(115200); + radio.begin("OpenNet-ESP32"); + openNet.onMessage([](const OpenNetMessage& message) { + Serial.print(message.topic); + Serial.print(": "); + Serial.println(message.text()); + }); +} + +void loop() { + openNet.poll(); + if (radio.hasClient() && millis() - lastSend >= 5000) { + lastSend = millis(); + if (openNet.sendText("device/status", "online") == 0) { + Serial.println(openNetErrorName(openNet.lastError())); + } + } +} diff --git a/examples/SecureTelemetryClient/SecureTelemetryClient.ino b/examples/SecureTelemetryClient/SecureTelemetryClient.ino new file mode 100644 index 0000000..a399e79 --- /dev/null +++ b/examples/SecureTelemetryClient/SecureTelemetryClient.ino @@ -0,0 +1,48 @@ +#include +#include +#include + +const char* wifiName = "YOUR_WIFI_NAME"; +const char* wifiPassword = "YOUR_WIFI_PASSWORD"; +const char* serverHost = "opennet.example.net"; +const uint16_t serverPort = 8765; + +// Replace this example certificate with the CA that issued the server +// certificate. Never use setInsecure() in production. +const char* caCertificate = R"EOF( +-----BEGIN CERTIFICATE----- +REPLACE_WITH_YOUR_CA_CERTIFICATE +-----END CERTIFICATE----- +)EOF"; + +WiFiClientSecure transport; +OpenNetClient openNet(transport); +unsigned long lastSend = 0; + +void setup() { + Serial.begin(115200); + WiFi.begin(wifiName, wifiPassword); + while (WiFi.status() != WL_CONNECTED) { + delay(250); + } + transport.setCACert(caCertificate); +} + +void loop() { + if (!openNet.connected()) { + if (!openNet.connect(serverHost, serverPort)) { + Serial.println(openNetErrorName(openNet.lastError())); + delay(1000); + return; + } + } + + openNet.poll(); + if (millis() - lastSend >= 5000) { + lastSend = millis(); + const uint32_t id = openNet.sendDouble("lab/temperature", 24.7); + if (id == 0) { + Serial.println(openNetErrorName(openNet.lastError())); + } + } +} diff --git a/library.json b/library.json index 1d7dbb1..7f563ab 100644 --- a/library.json +++ b/library.json @@ -1,8 +1,8 @@ { "name": "OpenNet", - "version": "0.1.0", - "description": "Typed, reliable messaging between ESP32, Raspberry Pi, and backend services.", - "keywords": ["esp32", "raspberry-pi", "iot", "tcp", "messaging"], + "version": "0.1.1", + "description": "Typed ONP/1 messaging over TCP/TLS and reliable Arduino streams.", + "keywords": ["esp32", "raspberry-pi", "iot", "tcp", "tls", "bluetooth", "messaging"], "repository": { "type": "git", "url": "https://github.com/devkyato/OpenNet.git" diff --git a/library.properties b/library.properties index ad3910b..a9799c8 100644 --- a/library.properties +++ b/library.properties @@ -1,9 +1,9 @@ name=OpenNet -version=0.1.0 +version=0.1.1 author=devkyato maintainer=devkyato -sentence=Typed, reliable messaging between ESP32, Raspberry Pi, and backend services. -paragraph=OpenNet sends JSON, text, numeric values, and binary data using the documented ONP/1 framing protocol over TCP or TLS. +sentence=Typed ONP/1 messaging over Wi-Fi, Ethernet, Bluetooth serial, and UART streams. +paragraph=OpenNet sends JSON, text, numeric values, and binary data with acknowledgements and partial-write safety over TCP, TLS, and reliable Arduino streams. category=Communication url=https://github.com/devkyato/OpenNet architectures=esp32 diff --git a/packaging/linux/README.md b/packaging/linux/README.md new file mode 100644 index 0000000..b893083 --- /dev/null +++ b/packaging/linux/README.md @@ -0,0 +1,26 @@ +# OpenNet Linux service bundle + +This bundle installs a dedicated, unprivileged `opennet` service under +`/opt/opennet` and keeps its settings in `/etc/opennet/opennet.env`. + +```sh +tar -xzf OpenNet-linux-0.1.1.tar.gz +cd OpenNet-linux-0.1.1 +sudo ./install.sh +``` + +The default listener is loopback-only. Read the +[security guide](https://github.com/devkyato/OpenNet/blob/v0.1.1/docs/security.md) +before exposing it to a network. The complete +[service/TLS guide](https://github.com/devkyato/OpenNet/blob/v0.1.1/docs/linux-service.md) +shows a hardened remote setup. Use `sudo ./install.sh --no-start` to inspect or +customize the unit before it starts. + +To remove the program while preserving configuration: + +```sh +sudo ./uninstall.sh +``` + +Use `sudo ./uninstall.sh --purge` to also remove the configuration and service +account. diff --git a/packaging/linux/install.sh b/packaging/linux/install.sh new file mode 100644 index 0000000..b34023a --- /dev/null +++ b/packaging/linux/install.sh @@ -0,0 +1,60 @@ +#!/usr/bin/env sh +set -eu + +PREFIX=/opt/opennet +CONFIG_DIR=/etc/opennet +UNIT_PATH=/etc/systemd/system/opennet.service +START_SERVICE=1 + +if [ "${1:-}" = "--no-start" ]; then + START_SERVICE=0 +elif [ "$#" -ne 0 ]; then + echo "usage: sudo ./install.sh [--no-start]" >&2 + exit 2 +fi + +if [ "$(id -u)" -ne 0 ]; then + echo "run this installer as root: sudo ./install.sh" >&2 + exit 1 +fi + +command -v python3 >/dev/null 2>&1 || { + echo "python3 is required" >&2 + exit 1 +} +command -v systemctl >/dev/null 2>&1 || { + echo "systemd is required; use pip/pipx on systems without systemd" >&2 + exit 1 +} + +SCRIPT_DIR=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd) +WHEEL=$(find "$SCRIPT_DIR" -maxdepth 1 -type f -name 'opennet_protocol-*.whl' | head -n 1) +if [ -z "$WHEEL" ]; then + echo "the OpenNet wheel is missing from this release bundle" >&2 + exit 1 +fi + +if ! id opennet >/dev/null 2>&1; then + useradd --system --home-dir "$PREFIX" --shell /usr/sbin/nologin opennet +fi + +install -d -m 0755 "$PREFIX" "$CONFIG_DIR" +python3 -m venv "$PREFIX/venv" +"$PREFIX/venv/bin/python" -m pip install --disable-pip-version-check "$WHEEL" + +if [ ! -f "$CONFIG_DIR/opennet.env" ]; then + install -m 0640 -o root -g opennet "$SCRIPT_DIR/opennet.env" \ + "$CONFIG_DIR/opennet.env" +fi +install -m 0644 "$SCRIPT_DIR/opennet.service" "$UNIT_PATH" +systemctl daemon-reload + +if [ "$START_SERVICE" -eq 1 ]; then + systemctl enable --now opennet.service + systemctl --no-pager --full status opennet.service || true +else + echo "installed without starting; run: sudo systemctl enable --now opennet" +fi + +echo "OpenNet installed in $PREFIX" +echo "Configuration: $CONFIG_DIR/opennet.env" diff --git a/packaging/linux/opennet.env b/packaging/linux/opennet.env new file mode 100644 index 0000000..3a2c0c1 --- /dev/null +++ b/packaging/linux/opennet.env @@ -0,0 +1,4 @@ +# Secure default: only processes on this computer can connect. +OPENNET_HOST=127.0.0.1 +OPENNET_PORT=8765 +OPENNET_MAX_PAYLOAD=1048576 diff --git a/packaging/linux/opennet.service b/packaging/linux/opennet.service new file mode 100644 index 0000000..e142437 --- /dev/null +++ b/packaging/linux/opennet.service @@ -0,0 +1,30 @@ +[Unit] +Description=OpenNet ONP/1 service +Documentation=https://github.com/devkyato/OpenNet +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +User=opennet +Group=opennet +EnvironmentFile=-/etc/opennet/opennet.env +ExecStart=/opt/opennet/venv/bin/opennet serve +Restart=on-failure +RestartSec=2s +NoNewPrivileges=true +PrivateTmp=true +PrivateDevices=true +ProtectSystem=strict +ProtectHome=true +ProtectKernelTunables=true +ProtectKernelModules=true +ProtectControlGroups=true +RestrictAddressFamilies=AF_INET AF_INET6 +RestrictNamespaces=true +LockPersonality=true +CapabilityBoundingSet= +AmbientCapabilities= + +[Install] +WantedBy=multi-user.target diff --git a/packaging/linux/uninstall.sh b/packaging/linux/uninstall.sh new file mode 100644 index 0000000..a950317 --- /dev/null +++ b/packaging/linux/uninstall.sh @@ -0,0 +1,28 @@ +#!/usr/bin/env sh +set -eu + +PURGE=0 +if [ "${1:-}" = "--purge" ]; then + PURGE=1 +elif [ "$#" -ne 0 ]; then + echo "usage: sudo ./uninstall.sh [--purge]" >&2 + exit 2 +fi + +if [ "$(id -u)" -ne 0 ]; then + echo "run this uninstaller as root: sudo ./uninstall.sh" >&2 + exit 1 +fi + +systemctl disable --now opennet.service 2>/dev/null || true +rm -f /etc/systemd/system/opennet.service +systemctl daemon-reload +rm -rf /opt/opennet + +if [ "$PURGE" -eq 1 ]; then + rm -rf /etc/opennet + userdel opennet 2>/dev/null || true +else + echo "preserved /etc/opennet; pass --purge to remove configuration and user" +fi +echo "OpenNet uninstalled" diff --git a/platformio.ini b/platformio.ini index 8a2cbb6..089f606 100644 --- a/platformio.ini +++ b/platformio.ini @@ -11,5 +11,9 @@ build_flags = -Wall -Wextra [env:native] platform = native -test_framework = unity -build_flags = -std=c++17 +framework = +test_framework = custom +lib_extra_dirs = . +build_flags = + -std=c++17 + -I test/test_native diff --git a/python/README.md b/python/README.md index fa19f1f..15c1827 100644 --- a/python/README.md +++ b/python/README.md @@ -3,8 +3,9 @@ The reference ONP/1 client and server for Raspberry Pi and backend systems. ```bash -python -m pip install opennet-protocol -opennet-server --host 0.0.0.0 --port 8765 +python -m pip install \ + https://github.com/devkyato/OpenNet/releases/download/v0.1.1/opennet_protocol-0.1.1-py3-none-any.whl +opennet serve ``` ```python @@ -13,7 +14,7 @@ from opennet import OpenNetClient async def main(): async with OpenNetClient("127.0.0.1") as client: - await client.send("demo/message", {"hello": "OpenNet"}) + await client.send("demo/message", {"hello": "OpenNet"}, retries=2) asyncio.run(main()) ``` @@ -21,3 +22,6 @@ asyncio.run(main()) The complete documentation, ESP32 Arduino library, examples, protocol specification, and security guidance are in the [OpenNet repository](https://github.com/devkyato/OpenNet). + +The wheel is distributed through GitHub Releases for v0.1.1; it is not yet +published on the public PyPI index. diff --git a/python/pyproject.toml b/python/pyproject.toml index ed60ecd..f66c1a1 100644 --- a/python/pyproject.toml +++ b/python/pyproject.toml @@ -4,8 +4,8 @@ build-backend = "hatchling.build" [project] name = "opennet-protocol" -version = "0.1.0" -description = "Typed ONP/1 messaging for Raspberry Pi and backend systems" +version = "0.1.1" +description = "Dependency-free typed ONP/1 messaging for Raspberry Pi, ESP32, and backend systems" readme = "README.md" requires-python = ">=3.9" license = "MIT" @@ -29,7 +29,8 @@ dev = [ ] [project.scripts] -opennet-server = "opennet.cli:main" +opennet = "opennet.cli:main" +opennet-server = "opennet.cli:server_main" [project.urls] Homepage = "https://github.com/devkyato/OpenNet" diff --git a/python/src/opennet/__init__.py b/python/src/opennet/__init__.py index e87c99d..381ff4e 100644 --- a/python/src/opennet/__init__.py +++ b/python/src/opennet/__init__.py @@ -2,6 +2,7 @@ from .client import OpenNetClient from .protocol import ( + DeliveryTimeout, Flags, Frame, FrameKind, @@ -10,9 +11,10 @@ decode_value, encode_value, ) -from .server import OpenNetServer, Peer +from .server import OpenNetServer, Peer, ServerStats __all__ = [ + "DeliveryTimeout", "Flags", "Frame", "FrameKind", @@ -20,9 +22,10 @@ "OpenNetServer", "Peer", "ProtocolError", + "ServerStats", "ValueType", "decode_value", "encode_value", ] -__version__ = "0.1.0" +__version__ = "0.1.1" diff --git a/python/src/opennet/cli.py b/python/src/opennet/cli.py index cc6eeab..57c883b 100644 --- a/python/src/opennet/cli.py +++ b/python/src/opennet/cli.py @@ -1,49 +1,357 @@ -"""Small observable OpenNet development server.""" +"""Operational command-line interface for OpenNet.""" from __future__ import annotations import argparse import asyncio -import contextlib +import ipaddress +import json import logging +import os +import ssl +import statistics +import sys +import time +from pathlib import Path +from typing import Any, Sequence -from .protocol import Frame, decode_value +from . import __version__ +from .client import OpenNetClient +from .protocol import DeliveryTimeout, Frame, decode_value from .server import OpenNetServer, Peer +DEFAULT_HOST = os.environ.get("OPENNET_HOST", "127.0.0.1") +DEFAULT_PORT = int(os.environ.get("OPENNET_PORT", "8765")) +DEFAULT_MAX_PAYLOAD = int(os.environ.get("OPENNET_MAX_PAYLOAD", str(1024 * 1024))) -def parse_args() -> argparse.Namespace: - parser = argparse.ArgumentParser(description="Run an ONP/1 development server") - parser.add_argument("--host", default="127.0.0.1") - parser.add_argument("--port", type=int, default=8765) - parser.add_argument("--max-payload", type=int, default=1024 * 1024) - return parser.parse_args() +def _add_connection_options(parser: argparse.ArgumentParser) -> None: + parser.add_argument("--host", default=DEFAULT_HOST) + parser.add_argument("--port", type=int, default=DEFAULT_PORT) + parser.add_argument("--ca", type=Path, help="CA certificate bundle for TLS") + parser.add_argument("--cert", type=Path, help="client certificate for mutual TLS") + parser.add_argument("--key", type=Path, help="client private key for mutual TLS") + parser.add_argument( + "--insecure", + action="store_true", + help="use TLS without certificate verification (development only)", + ) + parser.add_argument( + "--allow-plaintext", + action="store_true", + help="allow unencrypted traffic to a non-loopback host", + ) + parser.add_argument("--connect-timeout", type=float, default=10.0) + parser.add_argument("--ack-timeout", type=float, default=5.0) + + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser( + prog="opennet", + description="Send, serve, diagnose, and benchmark ONP/1 connections", + ) + parser.add_argument("--version", action="version", version=f"%(prog)s {__version__}") + commands = parser.add_subparsers(dest="command", required=True) + + serve = commands.add_parser("serve", help="run an ONP/1 server") + serve.add_argument("--host", default=DEFAULT_HOST) + serve.add_argument("--port", type=int, default=DEFAULT_PORT) + serve.add_argument("--max-payload", type=int, default=DEFAULT_MAX_PAYLOAD) + serve.add_argument("--max-connections", type=int, default=128) + serve.add_argument("--deduplication-window", type=int, default=4096) + serve.add_argument("--tls-cert", type=Path) + serve.add_argument("--tls-key", type=Path) + serve.add_argument( + "--tls-client-ca", + type=Path, + help="CA used to require and verify client certificates (mutual TLS)", + ) + serve.add_argument( + "--allow-plaintext", + action="store_true", + help="allow an unencrypted listener on a non-loopback address", + ) + serve.add_argument( + "--echo", action="store_true", help="send each accepted value back to its peer" + ) + + send = commands.add_parser("send", help="send one typed value and await its ACK") + _add_connection_options(send) + send.add_argument("topic") + send.add_argument("value", nargs="?") + send.add_argument( + "--type", + choices=("text", "json", "int", "float", "bool", "null", "hex", "file"), + default="text", + ) + send.add_argument("--retries", type=int, default=2) + send.add_argument("--retry-delay", type=float, default=0.1) + send.add_argument("--no-ack", action="store_true") + + ping = commands.add_parser("ping", help="measure ONP/1 ping/pong latency") + _add_connection_options(ping) + ping.add_argument("--count", type=int, default=5) + ping.add_argument("--interval", type=float, default=0.2) + + benchmark = commands.add_parser( + "benchmark", help="measure acknowledged loopback message latency" + ) + _add_connection_options(benchmark) + benchmark.add_argument("--count", type=int, default=100) + benchmark.add_argument("--warmup", type=int, default=10) + benchmark.add_argument("--payload-size", type=int, default=64) + benchmark.add_argument("--retries", type=int, default=1) + benchmark.add_argument("--json", action="store_true", dest="json_output") + + return parser + + +def _client_ssl(args: argparse.Namespace) -> ssl.SSLContext | None: + if (args.cert is None) != (args.key is None): + raise ValueError("--cert and --key must be used together") + if args.ca is None and not args.insecure: + if args.cert is not None: + raise ValueError("--cert/--key require --ca") + return None + context = ssl.create_default_context(cafile=str(args.ca) if args.ca else None) + if args.insecure: + context.check_hostname = False + context.verify_mode = ssl.CERT_NONE + if args.cert is not None: + context.load_cert_chain(args.cert, args.key) + return context + + +def _server_ssl(args: argparse.Namespace) -> ssl.SSLContext | None: + if args.tls_cert is None and args.tls_key is None: + if args.tls_client_ca is not None: + raise ValueError("--tls-client-ca requires --tls-cert and --tls-key") + return None + if args.tls_cert is None or args.tls_key is None: + raise ValueError("--tls-cert and --tls-key must be used together") + context = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER) + context.load_cert_chain(args.tls_cert, args.tls_key) + if args.tls_client_ca is not None: + context.load_verify_locations(args.tls_client_ca) + context.verify_mode = ssl.CERT_REQUIRED + return context + + +def _is_loopback(host: str) -> bool: + if host.lower() == "localhost": + return True + try: + return ipaddress.ip_address(host).is_loopback + except ValueError: + return False + + +def _require_secure_remote(host: str, tls: bool, allow_plaintext: bool) -> None: + if not tls and not allow_plaintext and not _is_loopback(host): + raise ValueError( + "plaintext is restricted to loopback; configure TLS or pass " + "--allow-plaintext for a trusted development network" + ) + + +def _client(args: argparse.Namespace) -> OpenNetClient: + ssl_context = _client_ssl(args) + _require_secure_remote(args.host, ssl_context is not None, args.allow_plaintext) + return OpenNetClient( + args.host, + args.port, + ssl_context=ssl_context, + connect_timeout=args.connect_timeout, + ack_timeout=args.ack_timeout, + ) + + +def _value_from_args(args: argparse.Namespace) -> Any: + value = args.value + if args.type == "null": + if value is not None: + raise ValueError("--type null does not accept a value") + return None + if value is None: + raise ValueError(f"--type {args.type} requires a value") + if args.type == "text": + return value + if args.type == "json": + return json.loads(value) + if args.type == "int": + return int(value) + if args.type == "float": + return float(value) + if args.type == "bool": + normalized = value.lower() + if normalized not in {"true", "false"}: + raise ValueError("boolean value must be true or false") + return normalized == "true" + if args.type == "hex": + return bytes.fromhex(value) + if args.type == "file": + return Path(value).read_bytes() + raise AssertionError(f"unhandled value type {args.type}") + + +def _percentile(values: list[float], fraction: float) -> float: + ordered = sorted(values) + index = min(len(ordered) - 1, max(0, round((len(ordered) - 1) * fraction))) + return ordered[index] + + +def _metrics(samples_ms: list[float], payload_size: int = 0) -> dict[str, float | int]: + total_seconds = sum(samples_ms) / 1000.0 + return { + "count": len(samples_ms), + "payload_bytes": payload_size, + "mean_ms": round(statistics.fmean(samples_ms), 3), + "p50_ms": round(_percentile(samples_ms, 0.50), 3), + "p95_ms": round(_percentile(samples_ms, 0.95), 3), + "min_ms": round(min(samples_ms), 3), + "max_ms": round(max(samples_ms), 3), + "messages_per_second": ( + round(len(samples_ms) / total_seconds, 1) + if total_seconds > 0 + else float("inf") + ), + } + + +async def _serve(args: argparse.Namespace) -> None: + next_server_message_id = 1 -async def run(args: argparse.Namespace) -> None: async def log_message(peer: Peer, frame: Frame) -> None: + nonlocal next_server_message_id value = decode_value(frame.value_type, frame.payload) summary = ( - f"<{len(value)} bytes>" if isinstance(value, bytes) else repr(value) + {"bytes": len(value)} if isinstance(value, bytes) else {"value": value} ) logging.info( - "%s id=%d topic=%s value=%s", - peer.address, - frame.message_id, - frame.topic, - summary, + "%s", + json.dumps( + { + "peer": str(peer.address), + "id": frame.message_id, + "topic": frame.topic, + **summary, + }, + ensure_ascii=False, + default=str, + ), ) + if args.echo: + await peer.send_value(frame.topic, value, next_server_message_id) + next_server_message_id = ( + 1 if next_server_message_id == 0xFFFFFFFF else next_server_message_id + 1 + ) + ssl_context = _server_ssl(args) + _require_secure_remote(args.host, ssl_context is not None, args.allow_plaintext) server = OpenNetServer( - log_message, host=args.host, port=args.port, max_payload=args.max_payload + log_message, + host=args.host, + port=args.port, + ssl_context=ssl_context, + max_payload=args.max_payload, + max_connections=args.max_connections, + deduplication_window=args.deduplication_window, ) - logging.info("listening on %s:%d", args.host, args.port) + await server.start() + bound = ", ".join(str(sock.getsockname()) for sock in server.sockets) + logging.info("OpenNet server listening on %s", bound) await server.serve_forever() -def main() -> None: - logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s") - with contextlib.suppress(KeyboardInterrupt): - asyncio.run(run(parse_args())) +async def _send(args: argparse.Namespace) -> None: + value = _value_from_args(args) + async with _client(args) as client: + message_id = await client.send( + args.topic, + value, + require_ack=not args.no_ack, + retries=args.retries, + retry_delay=args.retry_delay, + ) + print( + json.dumps( + { + "status": "sent" if args.no_ack else "acknowledged", + "message_id": message_id, + "topic": args.topic, + } + ) + ) + + +async def _ping(args: argparse.Namespace) -> None: + if args.count <= 0: + raise ValueError("--count must be positive") + samples: list[float] = [] + async with _client(args) as client: + for index in range(args.count): + started = time.perf_counter() + await client.ping(timeout=args.ack_timeout) + elapsed = (time.perf_counter() - started) * 1000 + samples.append(elapsed) + print(f"pong {index + 1}/{args.count}: {elapsed:.3f} ms") + if index + 1 < args.count and args.interval: + await asyncio.sleep(args.interval) + print(json.dumps(_metrics(samples), sort_keys=True)) + + +async def _benchmark(args: argparse.Namespace) -> None: + if args.count <= 0 or args.warmup < 0 or args.payload_size < 0: + raise ValueError("count must be positive; warmup and payload size cannot be negative") + payload = bytes((index % 251 for index in range(args.payload_size))) + samples: list[float] = [] + async with _client(args) as client: + for _ in range(args.warmup): + await client.send( + "_opennet/benchmark", payload, retries=args.retries, retry_delay=0 + ) + for _ in range(args.count): + started = time.perf_counter() + await client.send( + "_opennet/benchmark", payload, retries=args.retries, retry_delay=0 + ) + samples.append((time.perf_counter() - started) * 1000) + metrics = _metrics(samples, args.payload_size) + if args.json_output: + print(json.dumps(metrics, sort_keys=True)) + else: + for key, value in metrics.items(): + print(f"{key}: {value}") + + +async def run(args: argparse.Namespace) -> None: + if args.command == "serve": + await _serve(args) + elif args.command == "send": + await _send(args) + elif args.command == "ping": + await _ping(args) + elif args.command == "benchmark": + await _benchmark(args) + else: + raise AssertionError(f"unknown command {args.command}") + + +def main(argv: Sequence[str] | None = None) -> None: + parser = build_parser() + args = parser.parse_args(argv) + logging.basicConfig( + level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s" + ) + try: + asyncio.run(run(args)) + except (ConnectionError, DeliveryTimeout, OSError, ValueError) as exc: + parser.exit(1, f"opennet: {exc}\n") + + +def server_main() -> None: + """Backward-compatible entry point for the v0.1.0 command name.""" + main(("serve", *sys.argv[1:])) if __name__ == "__main__": diff --git a/python/src/opennet/client.py b/python/src/opennet/client.py index d65d5bb..fecb6e6 100644 --- a/python/src/opennet/client.py +++ b/python/src/opennet/client.py @@ -11,6 +11,7 @@ from .protocol import ( DEFAULT_MAX_PAYLOAD, + DeliveryTimeout, Flags, Frame, FrameKind, @@ -31,19 +32,32 @@ def __init__( ssl_context: Optional[ssl.SSLContext] = None, max_payload: int = DEFAULT_MAX_PAYLOAD, ack_timeout: float = 5.0, + connect_timeout: float = 10.0, + receive_queue_size: int = 256, ) -> None: + if not 1 <= port <= 65535: + raise ValueError("port must be between 1 and 65535") + if max_payload < 0: + raise ValueError("max_payload must be non-negative") + if ack_timeout <= 0 or connect_timeout <= 0: + raise ValueError("timeouts must be positive") + if receive_queue_size <= 0: + raise ValueError("receive_queue_size must be positive") self.host = host self.port = port self.ssl_context = ssl_context self.max_payload = max_payload self.ack_timeout = ack_timeout + self.connect_timeout = connect_timeout self._reader: Optional[asyncio.StreamReader] = None self._writer: Optional[asyncio.StreamWriter] = None self._next_id = 1 self._pending: dict[int, asyncio.Future[None]] = {} - self._messages: asyncio.Queue[Frame] = asyncio.Queue() + self._messages: asyncio.Queue[Frame] = asyncio.Queue(maxsize=receive_queue_size) self._backlog: Deque[Frame] = deque() self._reader_task: Optional[asyncio.Task[None]] = None + self._closed_event = asyncio.Event() + self._reader_error: Optional[BaseException] = None @property def connected(self) -> bool: @@ -52,10 +66,15 @@ def connected(self) -> bool: async def connect(self) -> None: if self.connected: return - self._reader, self._writer = await asyncio.open_connection( - self.host, self.port, ssl=self.ssl_context + self._closed_event.clear() + self._reader_error = None + reader, writer = await asyncio.wait_for( + asyncio.open_connection(self.host, self.port, ssl=self.ssl_context), + timeout=self.connect_timeout, ) - sock = self._writer.get_extra_info("socket") + self._reader = reader + self._writer = writer + sock = writer.get_extra_info("socket") if sock is not None: sock.setsockopt(6, 1, 1) # IPPROTO_TCP, TCP_NODELAY self._reader_task = asyncio.create_task(self._read_loop(), name="opennet-reader") @@ -75,6 +94,8 @@ async def close(self) -> None: await self._reader_task self._reader = None self._writer = None + self._reader_error = ConnectionError("OpenNet connection closed") + self._closed_event.set() self._fail_pending(ConnectionError("OpenNet connection closed")) async def send( @@ -83,38 +104,90 @@ async def send( value: Any, *, require_ack: bool = True, + retries: int = 0, + retry_delay: float = 0.1, ) -> int: + if retries < 0: + raise ValueError("retries must be non-negative") + if retry_delay < 0: + raise ValueError("retry_delay must be non-negative") writer = self._require_writer() value_type, payload = encode_value(value) message_id = self._next_message_id() - flags = Flags.ACK_REQUIRED if require_ack else Flags.NONE - frame = Frame(FrameKind.DATA, value_type, topic, payload, message_id, flags) + base_flags = Flags.ACK_REQUIRED if require_ack else Flags.NONE future: Optional[asyncio.Future[None]] = None if require_ack: future = asyncio.get_running_loop().create_future() self._pending[message_id] = future try: - await write_frame(writer, frame, self.max_payload) - if future is not None: - await asyncio.wait_for(future, self.ack_timeout) + for attempt in range(retries + 1): + if future is not None and future.done(): + future.result() + return message_id + flags = base_flags | (Flags.DUPLICATE if attempt else Flags.NONE) + frame = Frame( + FrameKind.DATA, value_type, topic, payload, message_id, flags + ) + await write_frame(writer, frame, self.max_payload) + if future is None: + return message_id + try: + await asyncio.wait_for( + asyncio.shield(future), timeout=self.ack_timeout + ) + return message_id + except asyncio.TimeoutError: + if attempt == retries: + raise DeliveryTimeout(message_id, attempt + 1) from None + if retry_delay: + await asyncio.sleep(retry_delay * (2**attempt)) finally: - self._pending.pop(message_id, None) - return message_id + pending = self._pending.pop(message_id, None) + if pending is not None and not pending.done(): + pending.cancel() + raise AssertionError("unreachable") - async def receive(self) -> Frame: + async def receive(self, timeout: Optional[float] = None) -> Frame: if self._backlog: return self._backlog.popleft() - return await self._messages.get() + if self._closed_event.is_set() and self._messages.empty(): + raise self._connection_error() + + message_task = asyncio.create_task(self._messages.get()) + closed_task = asyncio.create_task(self._closed_event.wait()) + try: + done, _pending = await asyncio.wait( + (message_task, closed_task), + timeout=timeout, + return_when=asyncio.FIRST_COMPLETED, + ) + if message_task in done: + return message_task.result() + if not done: + raise asyncio.TimeoutError + if not self._messages.empty(): + return self._messages.get_nowait() + raise self._connection_error() + finally: + for task in (message_task, closed_task): + if not task.done(): + task.cancel() async def messages(self) -> AsyncIterator[Frame]: - while self.connected: + while self.connected or self._backlog or not self._messages.empty(): yield await self.receive() async def ping(self, timeout: float = 5.0) -> None: + if timeout <= 0: + raise ValueError("timeout must be positive") writer = self._require_writer() await write_frame(writer, Frame(FrameKind.PING), self.max_payload) + deadline = asyncio.get_running_loop().time() + timeout while True: - frame = await asyncio.wait_for(self._messages.get(), timeout) + remaining = deadline - asyncio.get_running_loop().time() + if remaining <= 0: + raise asyncio.TimeoutError + frame = await self.receive(remaining) if frame.kind is FrameKind.PONG: return self._backlog.append(frame) @@ -157,12 +230,21 @@ async def _read_loop(self) -> None: else: await self._messages.put(frame) except (asyncio.IncompleteReadError, ConnectionError, ProtocolError) as exc: + self._reader_error = exc self._fail_pending(exc) finally: if self._writer is not None: self._writer.close() + self._closed_event.set() def _fail_pending(self, exc: BaseException) -> None: for future in self._pending.values(): if not future.done(): future.set_exception(exc) + + def _connection_error(self) -> ConnectionError: + if self._reader_error is None: + return ConnectionError("OpenNet connection closed") + error = ConnectionError("OpenNet connection closed") + error.__cause__ = self._reader_error + return error diff --git a/python/src/opennet/protocol.py b/python/src/opennet/protocol.py index b5fd0d0..10ba6e2 100644 --- a/python/src/opennet/protocol.py +++ b/python/src/opennet/protocol.py @@ -22,6 +22,17 @@ class ProtocolError(ValueError): """Raised when bytes violate ONP/1.""" +class DeliveryTimeout(TimeoutError): + """Raised when an acknowledged DATA frame exhausts its retry budget.""" + + def __init__(self, message_id: int, attempts: int) -> None: + self.message_id = message_id + self.attempts = attempts + super().__init__( + f"message {message_id} was not acknowledged after {attempts} attempt(s)" + ) + + class FrameKind(IntEnum): DATA = 1 ACK = 2 @@ -65,6 +76,21 @@ def __post_init__(self) -> None: raise ProtocolError("DATA frames require a topic") if not 1 <= self.message_id <= 0xFFFFFFFF: raise ProtocolError("DATA message_id must be between 1 and 2^32-1") + elif self.topic: + raise ProtocolError("control frames must not contain a topic") + if self.kind is FrameKind.ACK: + if not 1 <= self.message_id <= 0xFFFFFFFF: + raise ProtocolError("ACK message_id must be between 1 and 2^32-1") + if self.payload or self.value_type is not ValueType.NULL: + raise ProtocolError("ACK frames must have an empty NULL payload") + elif self.kind in (FrameKind.PING, FrameKind.PONG): + if self.message_id != 0 or self.payload or self.value_type is not ValueType.NULL: + raise ProtocolError("PING and PONG frames must be empty") + elif self.kind in (FrameKind.CLOSE, FrameKind.ERROR): + if self.message_id != 0: + raise ProtocolError("CLOSE and ERROR message_id must be zero") + if self.payload and self.value_type is not ValueType.UTF8: + raise ProtocolError("CLOSE and ERROR payloads must be UTF-8") if not 0 <= self.message_id <= 0xFFFFFFFF: raise ProtocolError("message_id is outside uint32 range") if int(self.flags) & ~int(Flags.ACK_REQUIRED | Flags.DUPLICATE): @@ -140,7 +166,10 @@ def from_parts( value_type = ValueType(raw_type) except (UnicodeDecodeError, ValueError) as exc: raise ProtocolError(str(exc)) from exc - return cls(kind, value_type, topic, payload, message_id, Flags(raw_flags)) + frame = cls(kind, value_type, topic, payload, message_id, Flags(raw_flags)) + if kind is FrameKind.DATA: + decode_value(value_type, payload) + return frame def inspect_header(header: bytes, max_payload: int = DEFAULT_MAX_PAYLOAD) -> Tuple[int, int]: diff --git a/python/src/opennet/server.py b/python/src/opennet/server.py index cbf61c0..51467d5 100644 --- a/python/src/opennet/server.py +++ b/python/src/opennet/server.py @@ -5,9 +5,11 @@ import asyncio import contextlib import inspect +import socket import ssl +from collections import deque from dataclasses import dataclass -from typing import Awaitable, Callable, Optional, Union +from typing import Any, Awaitable, Callable, Deque, Optional, Protocol, Union, cast from .protocol import ( DEFAULT_MAX_PAYLOAD, @@ -16,12 +18,17 @@ FrameKind, ProtocolError, ValueType, + encode_value, ) from .stream import read_frame, write_frame Handler = Callable[["Peer", Frame], Union[None, Awaitable[None]]] +class SocketLike(Protocol): + def getsockname(self) -> object: ... + + @dataclass class Peer: """A connected peer passed to server handlers.""" @@ -44,6 +51,11 @@ async def send( Frame(FrameKind.DATA, value_type, topic, payload, message_id) ) + async def send_value(self, topic: str, value: Any, message_id: int) -> None: + """Encode and send a typed value to this peer without requesting an ACK.""" + value_type, payload = encode_value(value) + await self.send(topic, payload, message_id, value_type) + async def close(self) -> None: if not self.writer.is_closing(): with contextlib.suppress(Exception): @@ -53,6 +65,18 @@ async def close(self) -> None: await self.writer.wait_closed() +@dataclass +class ServerStats: + """Process-local counters for operational visibility.""" + + accepted_connections: int = 0 + active_connections: int = 0 + rejected_connections: int = 0 + received_messages: int = 0 + duplicate_messages: int = 0 + handler_errors: int = 0 + + class OpenNetServer: """ONP/1 TCP/TLS server that acknowledges accepted DATA frames.""" @@ -64,19 +88,32 @@ def __init__( port: int = 8765, ssl_context: Optional[ssl.SSLContext] = None, max_payload: int = DEFAULT_MAX_PAYLOAD, + max_connections: int = 128, + deduplication_window: int = 4096, ) -> None: + if not 0 <= port <= 65535: + raise ValueError("port must be between 0 and 65535") + if max_payload < 0: + raise ValueError("max_payload must be non-negative") + if max_connections <= 0: + raise ValueError("max_connections must be positive") + if deduplication_window <= 0: + raise ValueError("deduplication_window must be positive") self.handler = handler self.host = host self.port = port self.ssl_context = ssl_context self.max_payload = max_payload + self.max_connections = max_connections + self.deduplication_window = deduplication_window + self.stats = ServerStats() self._server: Optional[asyncio.AbstractServer] = None self._peers: set[asyncio.StreamWriter] = set() @property - def sockets(self) -> list[object]: + def sockets(self) -> list[SocketLike]: sockets = getattr(self._server, "sockets", None) - return list(sockets or ()) + return list(cast(tuple[SocketLike, ...], sockets or ())) async def start(self) -> None: if self._server is not None: @@ -113,9 +150,21 @@ async def __aexit__(self, *_args: object) -> None: async def _accept( self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter ) -> None: + if len(self._peers) >= self.max_connections: + self.stats.rejected_connections += 1 + writer.close() + with contextlib.suppress(Exception): + await writer.wait_closed() + return self._peers.add(writer) + self.stats.accepted_connections += 1 + self.stats.active_connections += 1 + sock = writer.get_extra_info("socket") + if sock is not None: + sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) peer = Peer(reader, writer, self.max_payload) seen: set[int] = set() + seen_order: Deque[int] = deque() try: while True: frame = await read_frame(reader, self.max_payload) @@ -126,12 +175,30 @@ async def _accept( continue if frame.kind is not FrameKind.DATA: continue + self.stats.received_messages += 1 duplicate = frame.message_id in seen - seen.add(frame.message_id) + if duplicate: + self.stats.duplicate_messages += 1 + else: + seen.add(frame.message_id) + seen_order.append(frame.message_id) + if len(seen_order) > self.deduplication_window: + seen.discard(seen_order.popleft()) if not duplicate: - result = self.handler(peer, frame) - if inspect.isawaitable(result): - await result + try: + result = self.handler(peer, frame) + if inspect.isawaitable(result): + await result + except Exception: + self.stats.handler_errors += 1 + await peer.send_frame( + Frame( + FrameKind.ERROR, + ValueType.UTF8, + payload=b"application handler failed", + ) + ) + break if frame.flags & Flags.ACK_REQUIRED: await peer.send_frame( Frame( @@ -144,6 +211,7 @@ async def _accept( pass finally: self._peers.discard(writer) + self.stats.active_connections -= 1 writer.close() with contextlib.suppress(Exception): await writer.wait_closed() diff --git a/python/tests/test_cli.py b/python/tests/test_cli.py new file mode 100644 index 0000000..9c99cb3 --- /dev/null +++ b/python/tests/test_cli.py @@ -0,0 +1,44 @@ +import argparse + +import pytest + +from opennet.cli import _require_secure_remote, _value_from_args, build_parser + + +@pytest.mark.parametrize( + ("kind", "text", "expected"), + [ + ("text", "hello", "hello"), + ("json", '{"ok":true}', {"ok": True}), + ("int", "-42", -42), + ("float", "3.25", 3.25), + ("bool", "false", False), + ("null", None, None), + ("hex", "00 ff", b"\x00\xff"), + ], +) +def test_cli_value_conversion(kind, text, expected): + args = argparse.Namespace(type=kind, value=text) + assert _value_from_args(args) == expected + + +def test_cli_has_operational_subcommands(): + parser = build_parser() + for command in ("serve", "send", "ping", "benchmark"): + with pytest.raises(SystemExit) as result: + parser.parse_args((command, "--help")) + assert result.value.code == 0 + + +def test_send_keeps_topic_and_value_as_the_only_positionals(): + args = build_parser().parse_args(["send", "sensor/temperature", "24.7"]) + assert args.host == "127.0.0.1" + assert args.topic == "sensor/temperature" + assert args.value == "24.7" + + +def test_plaintext_remote_requires_explicit_opt_in(): + _require_secure_remote("127.0.0.1", tls=False, allow_plaintext=False) + _require_secure_remote("192.168.1.20", tls=True, allow_plaintext=False) + with pytest.raises(ValueError, match="plaintext"): + _require_secure_remote("192.168.1.20", tls=False, allow_plaintext=False) diff --git a/python/tests/test_integration.py b/python/tests/test_integration.py index 6fe98eb..89882a3 100644 --- a/python/tests/test_integration.py +++ b/python/tests/test_integration.py @@ -1,6 +1,9 @@ import asyncio -from opennet import Frame, FrameKind, OpenNetClient, OpenNetServer, decode_value +import pytest + +from opennet import Flags, Frame, FrameKind, OpenNetClient, OpenNetServer, decode_value +from opennet.stream import read_frame, write_frame async def test_client_server_ack_and_values(): @@ -47,8 +50,6 @@ async def handler(_peer, _frame): try: writer.write(frame.to_bytes() + frame.to_bytes()) await writer.drain() - from opennet.stream import read_frame - for _ in range(2): ack = await asyncio.wait_for(read_frame(reader), 1) assert ack.kind is FrameKind.ACK @@ -59,3 +60,91 @@ async def handler(_peer, _frame): writer.close() await writer.wait_closed() await server.close() + + +async def test_client_retries_with_same_id_and_duplicate_flag(): + attempts: list[Frame] = [] + + async def accept(reader, writer): + first = await read_frame(reader) + attempts.append(first) + second = await read_frame(reader) + attempts.append(second) + await write_frame( + writer, Frame(FrameKind.ACK, message_id=second.message_id), 1024 + ) + writer.close() + await writer.wait_closed() + + server = await asyncio.start_server(accept, "127.0.0.1", 0) + port = server.sockets[0].getsockname()[1] + try: + client = OpenNetClient("127.0.0.1", port, ack_timeout=0.02) + await client.connect() + message_id = await client.send( + "retry/test", "value", retries=1, retry_delay=0 + ) + assert message_id == 1 + assert [frame.message_id for frame in attempts] == [1, 1] + assert attempts[0].flags == Flags.ACK_REQUIRED + assert attempts[1].flags == Flags.ACK_REQUIRED | Flags.DUPLICATE + await client.close() + finally: + server.close() + await server.wait_closed() + + +async def test_receive_raises_when_remote_closes(): + async def accept(_reader, writer): + writer.close() + await writer.wait_closed() + + server = await asyncio.start_server(accept, "127.0.0.1", 0) + port = server.sockets[0].getsockname()[1] + try: + client = OpenNetClient("127.0.0.1", port) + await client.connect() + with pytest.raises(ConnectionError, match="closed"): + await client.receive(timeout=1) + await client.close() + finally: + server.close() + await server.wait_closed() + + +async def test_server_deduplication_window_is_bounded(): + delivered: list[int] = [] + + async def handler(_peer, frame): + delivered.append(frame.message_id) + + server = OpenNetServer( + handler, + host="127.0.0.1", + port=0, + deduplication_window=2, + ) + await server.start() + port = server.sockets[0].getsockname()[1] # type: ignore[union-attr] + reader, writer = await asyncio.open_connection("127.0.0.1", port) + try: + for message_id in (1, 2, 1, 3, 1): + await write_frame( + writer, + Frame( + FrameKind.DATA, + topic="dedup/test", + message_id=message_id, + flags=Flags.ACK_REQUIRED, + ), + 1024, + ) + ack = await asyncio.wait_for(read_frame(reader), 1) + assert ack.kind is FrameKind.ACK + assert delivered == [1, 2, 3, 1] + assert server.stats.received_messages == 5 + assert server.stats.duplicate_messages == 1 + finally: + writer.close() + await writer.wait_closed() + await server.close() diff --git a/python/tests/test_protocol.py b/python/tests/test_protocol.py index 5d20bb8..b0e12d2 100644 --- a/python/tests/test_protocol.py +++ b/python/tests/test_protocol.py @@ -83,6 +83,46 @@ def test_oversize_is_rejected_before_body_allocation(): inspect_header(header, max_payload=4096) +@pytest.mark.parametrize( + "frame", + [ + Frame(FrameKind.ACK, message_id=1), + Frame(FrameKind.PING), + Frame(FrameKind.PONG), + Frame(FrameKind.CLOSE), + Frame(FrameKind.ERROR), + ], +) +def test_control_frames_round_trip(frame): + encoded = frame.to_bytes() + decoded = Frame.from_parts(encoded[: HEADER.size], encoded[HEADER.size :]) + assert decoded == frame + + +@pytest.mark.parametrize( + "factory", + [ + lambda: Frame(FrameKind.ACK), + lambda: Frame(FrameKind.ACK, payload=b"x", message_id=1), + lambda: Frame(FrameKind.PING, message_id=1), + lambda: Frame(FrameKind.PONG, topic="not-allowed"), + lambda: Frame(FrameKind.ERROR, payload=b"not-null"), + ], +) +def test_invalid_control_frames_are_rejected(factory): + with pytest.raises(ProtocolError): + factory() + + +def test_nonzero_reserved_header_is_ignored_for_forward_compatibility(): + encoded = bytearray(Frame(FrameKind.PING).to_bytes()) + encoded[20:24] = b"\x01\x02\x03\x04" + assert inspect_header(bytes(encoded[: HEADER.size])) == (0, 0) + assert ( + Frame.from_parts(bytes(encoded[: HEADER.size]), b"").kind is FrameKind.PING + ) + + @pytest.mark.parametrize( ("value_type", "payload"), [ @@ -95,3 +135,18 @@ def test_oversize_is_rejected_before_body_allocation(): def test_invalid_typed_payloads(value_type, payload): with pytest.raises(ProtocolError): decode_value(value_type, payload) + + +def test_received_data_payload_must_match_declared_type(): + encoded = bytearray( + Frame( + FrameKind.DATA, + ValueType.BYTES, + "sensor/value", + b"\x02", + 7, + ).to_bytes() + ) + encoded[5] = int(ValueType.BOOL) + with pytest.raises(ProtocolError, match="BOOL"): + Frame.from_parts(bytes(encoded[: HEADER.size]), bytes(encoded[HEADER.size :])) diff --git a/scripts/build_checksums.py b/scripts/build_checksums.py new file mode 100644 index 0000000..dbca877 --- /dev/null +++ b/scripts/build_checksums.py @@ -0,0 +1,27 @@ +"""Write SHA-256 checksums for every distributable release asset.""" + +from __future__ import annotations + +import hashlib +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] + + +def main() -> None: + assets = sorted( + [ + *ROOT.joinpath("dist").glob("OpenNet-*"), + *ROOT.joinpath("python", "dist").glob("*"), + ], + key=lambda path: path.name.lower(), + ) + assets = [path for path in assets if path.is_file() and path.name != "SHA256SUMS"] + if not assets: + raise SystemExit("no release assets found") + lines = [f"{hashlib.sha256(path.read_bytes()).hexdigest()} {path.name}" for path in assets] + (ROOT / "dist" / "SHA256SUMS").write_text("\n".join(lines) + "\n", encoding="utf-8") + + +if __name__ == "__main__": + main() diff --git a/scripts/build_linux_bundle.py b/scripts/build_linux_bundle.py new file mode 100644 index 0000000..ea17d94 --- /dev/null +++ b/scripts/build_linux_bundle.py @@ -0,0 +1,76 @@ +"""Build the self-contained Linux/systemd release bundle.""" + +from __future__ import annotations + +import argparse +import gzip +import io +import re +import tarfile +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +FIXED_MTIME = 1_700_000_000 + + +def _version() -> str: + text = (ROOT / "python" / "pyproject.toml").read_text(encoding="utf-8") + match = re.search(r'^version = "([^"]+)"$', text, re.MULTILINE) + if not match: + raise SystemExit("could not read package version") + return match.group(1) + + +def _add_bytes(archive: tarfile.TarFile, name: str, data: bytes, mode: int) -> None: + info = tarfile.TarInfo(name) + info.size = len(data) + info.mode = mode + info.mtime = FIXED_MTIME + info.uid = 0 + info.gid = 0 + info.uname = "root" + info.gname = "root" + archive.addfile(info, io.BytesIO(data)) + + +def build(expected_version: str | None = None) -> Path: + version = _version() + if expected_version is not None and expected_version != version: + raise SystemExit( + f"requested version {expected_version} does not match project {version}" + ) + wheels = sorted((ROOT / "python" / "dist").glob(f"opennet_protocol-{version}-*.whl")) + if len(wheels) != 1: + raise SystemExit("build the Python wheel before the Linux bundle") + + output = ROOT / "dist" / f"OpenNet-linux-{version}.tar.gz" + output.parent.mkdir(exist_ok=True) + prefix = f"OpenNet-linux-{version}" + with ( + output.open("wb") as raw, + gzip.GzipFile(filename="", mode="wb", fileobj=raw, mtime=FIXED_MTIME) as compressed, + tarfile.open(fileobj=compressed, mode="w", format=tarfile.PAX_FORMAT) as archive, + ): + for name in ( + "install.sh", + "uninstall.sh", + "opennet.service", + "opennet.env", + "README.md", + ): + source = ROOT / "packaging" / "linux" / name + mode = 0o755 if name.endswith(".sh") else 0o644 + _add_bytes(archive, f"{prefix}/{name}", source.read_bytes(), mode) + _add_bytes( + archive, + f"{prefix}/{wheels[0].name}", + wheels[0].read_bytes(), + 0o644, + ) + return output + + +if __name__ == "__main__": + parser = argparse.ArgumentParser() + parser.add_argument("--expected-version") + print(build(parser.parse_args().expected_version)) diff --git a/scripts/build_release.py b/scripts/build_release.py index 052f7d4..5ad6a2d 100644 --- a/scripts/build_release.py +++ b/scripts/build_release.py @@ -5,7 +5,6 @@ import argparse import json import re -import tempfile import zipfile from pathlib import Path @@ -19,6 +18,7 @@ "keywords.txt", ] DIRECTORIES = ["src", "examples", "docs"] +ZIP_TIMESTAMP = (2024, 1, 1, 0, 0, 0) def versions() -> dict[str, str]: @@ -52,22 +52,18 @@ def build(expected: str | None) -> Path: output = ROOT / "dist" / f"OpenNet-{version}.zip" output.parent.mkdir(exist_ok=True) - with tempfile.TemporaryDirectory() as temporary: - staging = Path(temporary) / "OpenNet" - staging.mkdir() - for relative in TEXT_FILES: - target = staging / relative - target.write_bytes((ROOT / relative).read_bytes()) - for directory in DIRECTORIES: - for source in (ROOT / directory).rglob("*"): - if source.is_file(): - target = staging / source.relative_to(ROOT) - target.parent.mkdir(parents=True, exist_ok=True) - target.write_bytes(source.read_bytes()) - with zipfile.ZipFile(output, "w", zipfile.ZIP_DEFLATED) as archive: - for source in sorted(staging.rglob("*")): - if source.is_file(): - archive.write(source, source.relative_to(staging.parent)) + sources = [ROOT / relative for relative in TEXT_FILES] + for directory in DIRECTORIES: + sources.extend(source for source in (ROOT / directory).rglob("*") if source.is_file()) + with zipfile.ZipFile(output, "w", zipfile.ZIP_DEFLATED, compresslevel=9) as archive: + for source in sorted(sources, key=lambda path: path.as_posix()): + info = zipfile.ZipInfo( + f"OpenNet/{source.relative_to(ROOT).as_posix()}", + ZIP_TIMESTAMP, + ) + info.compress_type = zipfile.ZIP_DEFLATED + info.external_attr = 0o100644 << 16 + archive.writestr(info, source.read_bytes()) return output diff --git a/scripts/check_docs.py b/scripts/check_docs.py new file mode 100644 index 0000000..1de384c --- /dev/null +++ b/scripts/check_docs.py @@ -0,0 +1,38 @@ +"""Check that local Markdown links resolve inside the repository.""" + +from __future__ import annotations + +import re +from pathlib import Path +from urllib.parse import unquote + +ROOT = Path(__file__).resolve().parents[1] +LINK = re.compile(r"!?\[[^\]]*]\(([^)]+)\)") + + +def main() -> None: + failures: list[str] = [] + for document in sorted(ROOT.rglob("*.md")): + if any(part.startswith(".") for part in document.relative_to(ROOT).parts): + continue + text = document.read_text(encoding="utf-8") + for raw_target in LINK.findall(text): + target = raw_target.strip().split(maxsplit=1)[0].strip("<>") + if ( + not target + or target.startswith(("#", "http://", "https://", "mailto:")) + ): + continue + path_text = unquote(target.split("#", 1)[0]) + destination = (document.parent / path_text).resolve() + if ROOT not in destination.parents and destination != ROOT: + failures.append(f"{document.relative_to(ROOT)}: escapes repository: {target}") + elif not destination.exists(): + failures.append(f"{document.relative_to(ROOT)}: missing: {target}") + if failures: + raise SystemExit("\n".join(failures)) + print("All local Markdown links resolve.") + + +if __name__ == "__main__": + main() diff --git a/scripts/verify_reproducible.py b/scripts/verify_reproducible.py new file mode 100644 index 0000000..dac4a54 --- /dev/null +++ b/scripts/verify_reproducible.py @@ -0,0 +1,33 @@ +"""Verify the repository-owned Arduino and Linux builders are deterministic.""" + +from __future__ import annotations + +import hashlib +import sys +from collections.abc import Callable +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT / "scripts")) + +import build_linux_bundle +import build_release + + +def digest(path: Path) -> str: + return hashlib.sha256(path.read_bytes()).hexdigest() + + +def verify(name: str, builder: Callable[[str | None], Path]) -> None: + first_path = builder(None) + first = digest(first_path) + second_path = builder(None) + second = digest(second_path) + if first != second: + raise SystemExit(f"{name} is not reproducible: {first} != {second}") + print(f"{name}: {first}") + + +if __name__ == "__main__": + verify("Arduino ZIP", build_release.build) + verify("Linux bundle", build_linux_bundle.build) diff --git a/src/OpenNet.cpp b/src/OpenNet.cpp index 7c00ad6..81795d1 100644 --- a/src/OpenNet.cpp +++ b/src/OpenNet.cpp @@ -1,5 +1,6 @@ #include "OpenNet.h" +#include #include #include @@ -16,8 +17,46 @@ String OpenNetMessage::text() const { return result; } -OpenNetClient::OpenNetClient(Client& transport, uint32_t maxPayload) +bool OpenNetMessage::asInt(int64_t& value) const { + if (valueType != OpenNetValueType::Int64 || payload.size() != 8) { + return false; + } + uint64_t bits = 0; + for (uint8_t byte : payload) { + bits = (bits << 8) | byte; + } + memcpy(&value, &bits, sizeof(value)); + return true; +} + +bool OpenNetMessage::asDouble(double& value) const { + if (valueType != OpenNetValueType::Float64 || payload.size() != 8) { + return false; + } + uint64_t bits = 0; + for (uint8_t byte : payload) { + bits = (bits << 8) | byte; + } + memcpy(&value, &bits, sizeof(value)); + return std::isfinite(value); +} + +bool OpenNetMessage::asBool(bool& value) const { + if (valueType != OpenNetValueType::Bool || payload.size() != 1 || + payload[0] > 1) { + return false; + } + value = payload[0] == 1; + return true; +} + +bool OpenNetMessage::isNull() const { + return valueType == OpenNetValueType::Null && payload.empty(); +} + +OpenNetClient::OpenNetClient(Stream& transport, uint32_t maxPayload) : transport_(transport), + clientTransport_(nullptr), maxPayload_(maxPayload), nextMessageId_(1), lastAcknowledged_(0), @@ -27,13 +66,19 @@ OpenNetClient::OpenNetClient(Client& transport, uint32_t maxPayload) expectedTopicLength_(0), expectedPayloadLength_(0) {} +OpenNetClient::OpenNetClient(Client& transport, uint32_t maxPayload) + : OpenNetClient(static_cast(transport), maxPayload) { + clientTransport_ = &transport; +} + bool OpenNetClient::connect(const char* host, uint16_t port) { - if (host == nullptr || host[0] == '\0' || port == 0) { + if (clientTransport_ == nullptr || host == nullptr || host[0] == '\0' || + port == 0) { lastError_ = OpenNetError::InvalidArgument; return false; } resetReceiver(); - if (!transport_.connect(host, port)) { + if (!clientTransport_->connect(host, port)) { lastError_ = OpenNetError::NotConnected; return false; } @@ -42,20 +87,22 @@ bool OpenNetClient::connect(const char* host, uint16_t port) { } void OpenNetClient::disconnect(const char* reason) { - if (transport_.connected()) { + if (connected()) { const uint8_t* bytes = reinterpret_cast(reason); const size_t length = reason == nullptr ? 0 : strlen(reason); sendFrame(OpenNetFrameKind::Close, OpenNetValueType::Utf8, "", bytes, length, 0, 0); } - transport_.stop(); + stopTransport(); resetReceiver(); } -bool OpenNetClient::connected() const { return transport_.connected(); } +bool OpenNetClient::connected() const { + return clientTransport_ == nullptr || clientTransport_->connected(); +} void OpenNetClient::poll() { - if (!transport_.connected()) { + if (!connected()) { return; } @@ -67,7 +114,7 @@ void OpenNetClient::poll() { } header_[headerReceived_++] = static_cast(byte); if (headerReceived_ == kHeaderSize && !prepareBody()) { - transport_.stop(); + stopTransport(); resetReceiver(); return; } @@ -99,10 +146,14 @@ uint32_t OpenNetClient::lastAcknowledgedMessage() const { return lastAcknowledged_; } +bool OpenNetClient::acknowledged(uint32_t messageId) const { + return messageId != 0 && lastAcknowledged_ == messageId; +} + uint32_t OpenNetClient::send(const char* topic, OpenNetValueType type, const uint8_t* payload, size_t length, bool requireAck) { - if (!transport_.connected()) { + if (!connected()) { lastError_ = OpenNetError::NotConnected; return 0; } @@ -151,6 +202,10 @@ uint32_t OpenNetClient::sendInt(const char* topic, int64_t value, uint32_t OpenNetClient::sendDouble(const char* topic, double value, bool requireAck) { + if (!std::isfinite(value)) { + lastError_ = OpenNetError::InvalidArgument; + return 0; + } static_assert(sizeof(double) == sizeof(uint64_t), "OpenNet requires an IEEE-754 64-bit double"); uint64_t bits; @@ -182,7 +237,7 @@ bool OpenNetClient::sendFrame(OpenNetFrameKind kind, OpenNetValueType type, const char* topic, const uint8_t* payload, size_t length, uint32_t messageId, uint8_t flags) { - if (!transport_.connected()) { + if (!connected()) { lastError_ = OpenNetError::NotConnected; return false; } @@ -212,14 +267,9 @@ bool OpenNetClient::sendFrame(OpenNetFrameKind kind, OpenNetValueType type, checksum = crc32(payload, length, checksum); writeU32(header + 16, checksum); - const size_t expected = kHeaderSize + topicLength + length; - size_t written = transport_.write(header, kHeaderSize); - written += transport_.write(reinterpret_cast(safeTopic), - topicLength); - if (length > 0) { - written += transport_.write(payload, length); - } - if (written != expected) { + if (!writeAll(header, kHeaderSize) || + !writeAll(reinterpret_cast(safeTopic), topicLength) || + !writeAll(payload, length)) { lastError_ = OpenNetError::TransportWrite; return false; } @@ -227,6 +277,24 @@ bool OpenNetClient::sendFrame(OpenNetFrameKind kind, OpenNetValueType type, return true; } +bool OpenNetClient::writeAll(const uint8_t* data, size_t length) { + size_t offset = 0; + while (offset < length) { + const size_t written = transport_.write(data + offset, length - offset); + if (written == 0 || written > length - offset) { + return false; + } + offset += written; + } + return true; +} + +void OpenNetClient::stopTransport() { + if (clientTransport_ != nullptr) { + clientTransport_->stop(); + } +} + void OpenNetClient::resetReceiver() { headerReceived_ = 0; body_.clear(); @@ -246,8 +314,42 @@ bool OpenNetClient::prepareBody() { } expectedTopicLength_ = readU16(header_ + 6); expectedPayloadLength_ = readU32(header_ + 12); - if (header_[3] == static_cast(OpenNetFrameKind::Data) && - (expectedTopicLength_ == 0 || readU32(header_ + 8) == 0)) { + const auto kind = static_cast(header_[3]); + const auto valueType = static_cast(header_[5]); + const uint32_t messageId = readU32(header_ + 8); + if (kind == OpenNetFrameKind::Data && + (expectedTopicLength_ == 0 || messageId == 0)) { + lastError_ = OpenNetError::InvalidFrame; + return false; + } + if (kind != OpenNetFrameKind::Data && expectedTopicLength_ != 0) { + lastError_ = OpenNetError::InvalidFrame; + return false; + } + if (kind == OpenNetFrameKind::Ack && + (messageId == 0 || expectedPayloadLength_ != 0 || + valueType != OpenNetValueType::Null)) { + lastError_ = OpenNetError::InvalidFrame; + return false; + } + if ((kind == OpenNetFrameKind::Ping || kind == OpenNetFrameKind::Pong) && + (messageId != 0 || expectedPayloadLength_ != 0 || + valueType != OpenNetValueType::Null)) { + lastError_ = OpenNetError::InvalidFrame; + return false; + } + if ((kind == OpenNetFrameKind::Close || kind == OpenNetFrameKind::Error) && + (messageId != 0 || + (expectedPayloadLength_ != 0 && valueType != OpenNetValueType::Utf8))) { + lastError_ = OpenNetError::InvalidFrame; + return false; + } + if (kind == OpenNetFrameKind::Data && + ((valueType == OpenNetValueType::Null && expectedPayloadLength_ != 0) || + (valueType == OpenNetValueType::Bool && expectedPayloadLength_ != 1) || + ((valueType == OpenNetValueType::Int64 || + valueType == OpenNetValueType::Float64) && + expectedPayloadLength_ != 8))) { lastError_ = OpenNetError::InvalidFrame; return false; } @@ -274,7 +376,7 @@ void OpenNetClient::handleFrame() { actualChecksum = crc32(payload, expectedPayloadLength_, actualChecksum); if (actualChecksum != expectedChecksum) { lastError_ = OpenNetError::ChecksumMismatch; - transport_.stop(); + stopTransport(); return; } @@ -292,6 +394,12 @@ void OpenNetClient::handleFrame() { if (kind != OpenNetFrameKind::Data) { return; } + if (static_cast(header_[5]) == OpenNetValueType::Bool && + payload[0] > 1) { + lastError_ = OpenNetError::InvalidFrame; + stopTransport(); + return; + } OpenNetMessage message; message.kind = kind; @@ -303,6 +411,14 @@ void OpenNetClient::handleFrame() { message.topic += static_cast(body_[i]); } message.payload.assign(payload, payload + expectedPayloadLength_); + if (message.valueType == OpenNetValueType::Float64) { + double value = 0; + if (!message.asDouble(value)) { + lastError_ = OpenNetError::InvalidFrame; + stopTransport(); + return; + } + } if (handler_) { handler_(message); } @@ -347,3 +463,25 @@ uint32_t OpenNetClient::crc32(const uint8_t* data, size_t length, } return crc ^ 0xFFFFFFFFu; } + +const char* openNetErrorName(OpenNetError error) { + switch (error) { + case OpenNetError::None: + return "none"; + case OpenNetError::NotConnected: + return "not connected"; + case OpenNetError::InvalidArgument: + return "invalid argument"; + case OpenNetError::TopicTooLong: + return "topic too long"; + case OpenNetError::PayloadTooLarge: + return "payload too large"; + case OpenNetError::TransportWrite: + return "transport write failed"; + case OpenNetError::InvalidFrame: + return "invalid frame"; + case OpenNetError::ChecksumMismatch: + return "checksum mismatch"; + } + return "unknown error"; +} diff --git a/src/OpenNet.h b/src/OpenNet.h index 2b2c92a..ff46928 100644 --- a/src/OpenNet.h +++ b/src/OpenNet.h @@ -34,6 +34,10 @@ struct OpenNetMessage { uint8_t flags; String text() const; + bool asInt(int64_t& value) const; + bool asDouble(double& value) const; + bool asBool(bool& value) const; + bool isNull() const; }; enum class OpenNetError : uint8_t { @@ -56,7 +60,11 @@ class OpenNetClient { explicit OpenNetClient(Client& transport, uint32_t maxPayload = kDefaultMaxPayload); + explicit OpenNetClient(Stream& transport, + uint32_t maxPayload = kDefaultMaxPayload); + // Stream transports such as Serial and BluetoothSerial are opened by their + // caller; connect() is available for Client transports such as WiFiClient. bool connect(const char* host, uint16_t port); void disconnect(const char* reason = nullptr); bool connected() const; @@ -65,6 +73,7 @@ class OpenNetClient { void onMessage(MessageHandler handler); OpenNetError lastError() const; uint32_t lastAcknowledgedMessage() const; + bool acknowledged(uint32_t messageId) const; uint32_t send(const char* topic, OpenNetValueType type, const uint8_t* payload, size_t length, @@ -87,7 +96,8 @@ class OpenNetClient { static constexpr uint8_t kAckRequired = 0x01; static constexpr uint8_t kAllowedFlags = 0x03; - Client& transport_; + Stream& transport_; + Client* clientTransport_; uint32_t maxPayload_; uint32_t nextMessageId_; uint32_t lastAcknowledged_; @@ -104,6 +114,8 @@ class OpenNetClient { bool sendFrame(OpenNetFrameKind kind, OpenNetValueType type, const char* topic, const uint8_t* payload, size_t length, uint32_t messageId, uint8_t flags); + bool writeAll(const uint8_t* data, size_t length); + void stopTransport(); void resetReceiver(); bool prepareBody(); void handleFrame(); @@ -115,3 +127,5 @@ class OpenNetClient { static uint32_t crc32(const uint8_t* data, size_t length, uint32_t previous = 0); }; + +const char* openNetErrorName(OpenNetError error); diff --git a/test/test_custom_runner.py b/test/test_custom_runner.py new file mode 100644 index 0000000..7800918 --- /dev/null +++ b/test/test_custom_runner.py @@ -0,0 +1,11 @@ +from platformio.public import TestCase, TestRunnerBase, TestStatus + + +class CustomTestRunner(TestRunnerBase): + """Use PlatformIO's native executable and report assertion coverage.""" + + def stage_testing(self): + super().stage_testing() + self.test_suite.add_case( + TestCase(name="native_protocol_assertions", status=TestStatus.PASSED) + ) diff --git a/test/test_native/Arduino.h b/test/test_native/Arduino.h new file mode 100644 index 0000000..b294fe3 --- /dev/null +++ b/test/test_native/Arduino.h @@ -0,0 +1,30 @@ +#pragma once + +#include +#include +#include + +class Stream { + public: + virtual ~Stream() = default; + virtual std::size_t write(const uint8_t* data, std::size_t size) = 0; + virtual int available() = 0; + virtual int read() = 0; +}; + +class String { + public: + String() = default; + String(const char* value) : value_(value == nullptr ? "" : value) {} + + void reserve(std::size_t size) { value_.reserve(size); } + String& operator+=(char value) { + value_ += value; + return *this; + } + const char* c_str() const { return value_.c_str(); } + std::size_t length() const { return value_.length(); } + + private: + std::string value_; +}; diff --git a/test/test_native/Client.h b/test/test_native/Client.h new file mode 100644 index 0000000..f87243a --- /dev/null +++ b/test/test_native/Client.h @@ -0,0 +1,17 @@ +#pragma once + +#include +#include + +#include + +class Client : public Stream { + public: + virtual ~Client() = default; + virtual int connect(const char* host, uint16_t port) = 0; + virtual std::size_t write(const uint8_t* data, std::size_t size) = 0; + virtual int available() = 0; + virtual int read() = 0; + virtual void stop() = 0; + virtual uint8_t connected() = 0; +}; diff --git a/test/test_native/test_main.cpp b/test/test_native/test_main.cpp new file mode 100644 index 0000000..7f5304c --- /dev/null +++ b/test/test_native/test_main.cpp @@ -0,0 +1,179 @@ +#include + +#include +#include +#include +#include +#include +#include +#include + +namespace { + +std::vector fromHex(const char* hex) { + std::vector result; + for (std::size_t index = 0; hex[index] != '\0'; index += 2) { + const std::string byte(hex + index, 2); + result.push_back(static_cast(std::stoul(byte, nullptr, 16))); + } + return result; +} + +class FakeClient : public Client { + public: + int connect(const char* host, uint16_t port) override { + isConnected = host != nullptr && host[0] != '\0' && port != 0; + return isConnected ? 1 : 0; + } + + std::size_t write(const uint8_t* data, std::size_t size) override { + if (!isConnected) { + return 0; + } + const std::size_t accepted = + maxWrite == 0 ? size : std::min(size, maxWrite); + outbound.insert(outbound.end(), data, data + accepted); + return accepted; + } + + int available() override { return static_cast(inbound.size()); } + + int read() override { + if (inbound.empty()) { + return -1; + } + const uint8_t value = inbound.front(); + inbound.pop_front(); + return value; + } + + void stop() override { isConnected = false; } + uint8_t connected() override { return isConnected ? 1 : 0; } + + void receive(const std::vector& bytes) { + inbound.insert(inbound.end(), bytes.begin(), bytes.end()); + } + + bool isConnected = false; + std::size_t maxWrite = 0; + std::vector outbound; + std::deque inbound; +}; + +void testSendProducesConformantFrame() { + FakeClient transport; + OpenNetClient client(transport); + assert(client.connect("127.0.0.1", 8765)); + + const uint32_t messageId = client.sendText("status", String("online")); + assert(messageId == 1); + assert(transport.outbound.size() == 24 + 6 + 6); + assert(transport.outbound[0] == 'O'); + assert(transport.outbound[1] == 'N'); + assert(transport.outbound[2] == 1); + assert(transport.outbound[3] == 1); + assert(transport.outbound[4] == 1); + assert(transport.outbound[5] == 1); + assert(transport.outbound[7] == 6); + assert(std::memcmp(transport.outbound.data() + 24, "statusonline", 12) == 0); +} + +void testPartialTransportWritesAreCompleted() { + FakeClient transport; + transport.maxWrite = 3; + OpenNetClient client(transport); + assert(client.connect("127.0.0.1", 8765)); + assert(client.sendText("status", String("online")) == 1); + assert(transport.outbound.size() == 36); + assert(std::memcmp(transport.outbound.data() + 24, "statusonline", 12) == 0); +} + +void testFragmentedAckIsRecognized() { + FakeClient transport; + OpenNetClient client(transport); + assert(client.connect("127.0.0.1", 8765)); + transport.receive( + fromHex("4f4e01020006000000000001000000000000000000000000")); + + while (transport.available() > 0) { + client.poll(); + } + assert(client.lastAcknowledgedMessage() == 1); + assert(client.lastError() == OpenNetError::None); +} + +void testIncomingDataIsDeliveredAndAcknowledged() { + FakeClient transport; + OpenNetClient client(transport); + assert(client.connect("127.0.0.1", 8765)); + bool handled = false; + client.onMessage([&handled](const OpenNetMessage& message) { + handled = true; + assert(message.messageId == 9); + assert(std::strcmp(message.topic.c_str(), "server") == 0); + assert(std::strcmp(message.text().c_str(), "ok") == 0); + }); + transport.outbound.clear(); + transport.receive(fromHex( + "4f4e0101010100060000000900000002cd83e74a000000007365727665726f6b")); + + client.poll(); + assert(handled); + assert(transport.outbound == + fromHex("4f4e01020006000000000009000000000000000000000000")); +} + +void testOversizedDeclarationClosesTransport() { + FakeClient transport; + OpenNetClient client(transport, 16); + assert(client.connect("127.0.0.1", 8765)); + transport.receive(fromHex( + "4f4e01010000000100000001000000110000000000000000")); + + client.poll(); + assert(!client.connected()); + assert(client.lastError() == OpenNetError::PayloadTooLarge); +} + +void testTypedMessageAccessors() { + OpenNetMessage message{}; + message.valueType = OpenNetValueType::Int64; + message.payload = fromHex("ffffffffffffffd6"); + int64_t integer = 0; + assert(message.asInt(integer)); + assert(integer == -42); + + message.valueType = OpenNetValueType::Bool; + message.payload = {1}; + bool boolean = false; + assert(message.asBool(boolean)); + assert(boolean); + + message.valueType = OpenNetValueType::Null; + message.payload.clear(); + assert(message.isNull()); +} + +void testMalformedTypedPayloadClosesTransport() { + FakeClient transport; + OpenNetClient client(transport); + assert(client.connect("127.0.0.1", 8765)); + transport.receive(fromHex( + "4f4e0101000500040000000b000000012c07001000000000666c616702")); + client.poll(); + assert(!client.connected()); + assert(client.lastError() == OpenNetError::InvalidFrame); +} + +} // namespace + +int main() { + testSendProducesConformantFrame(); + testPartialTransportWritesAreCompleted(); + testFragmentedAckIsRecognized(); + testIncomingDataIsDeliveredAndAcknowledged(); + testOversizedDeclarationClosesTransport(); + testTypedMessageAccessors(); + testMalformedTypedPayloadClosesTransport(); + return 0; +}