diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2bcf02b..374b97c 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -44,6 +44,12 @@ jobs: poetry run python scripts/gen_proto.py git diff --exit-code || (echo "Error: Generated protobuf stubs are out of sync with pzem_004t.proto. Run scripts/gen_proto.py and commit." && exit 1) + - name: Run Ruff linter + run: poetry run ruff check . + + - name: Run Ruff formatter check + run: poetry run ruff format --check . + - name: Run strict type checking (Pyrefly) run: poetry run pyrefly check diff --git a/.vscode/settings.json b/.vscode/settings.json index 6a742d9..fd15719 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -1,4 +1,12 @@ { "python-envs.defaultEnvManager": "ms-python.python:poetry", - "python-envs.defaultPackageManager": "ms-python.python:poetry" + "python-envs.defaultPackageManager": "ms-python.python:poetry", + "[python]": { + "editor.defaultFormatter": "charliermarsh.ruff", + "editor.formatOnSave": true, + "editor.codeActionsOnSave": { + "source.fixAll": "explicit", + "source.organizeImports": "explicit" + } + } } \ No newline at end of file diff --git a/README.md b/README.md index e377e7f..af796c5 100644 --- a/README.md +++ b/README.md @@ -11,42 +11,44 @@ interact with the gRPC microservice. ```mermaid flowchart LR - subgraph RESTConsumer["External REST Consumer"] - Web["Web / Mobile App / cURL"] + subgraph Microcontrollers["IoT Microcontrollers (e.g. ESP32)"] + ESP32["ESP32 + PZEM-004T\n(apps/mqtt_sensor_node)\n[paho-mqtt]"] end - subgraph Gateway["REST-to-gRPC Gateway (FastAPI)"] - GW["gateway/ (app.py)"] + subgraph MQTTBroker["MQTT Messaging Broker (Port 1883)"] + Broker["Eclipse Mosquitto\n(devices/+/telemetry)"] end - subgraph Device["PZEM-004t Device Gateway (client)"] - A1["client/ (telemetry.py)"] - A2["device/ (pzem_004t.py)"] + subgraph Bridges["Ingestion Bridges"] + Bridge["MQTT-to-gRPC Bridge\n(apps/mqtt_bridge)\n[paho-mqtt + grpc.aio]"] end - subgraph Collector["Telemetry Collector Service (server)"] - B1["server/ (servicer.py + app.py)"] - B2["Health Checking (grpc.health.v1)"] - B3["Interceptors (Tracing, Metrics, Recovery)"] - B4["in-memory store + logging"] + subgraph Collector["Telemetry Collector Service (Port 50051)"] + Coll["Collector Service\n(apps/collector)\n(gRPC Ingestion & Pub/Sub Hub)"] end - Web --> |"HTTP/JSON REST API (port 8000)"| GW - GW --> |"gRPC Unary / Client-Stream"| B1 - A1 <==> |"gRPC (HTTP/2) - 4 RPC call types (port 50051)"| B1 - A1 -- "simulated readings" --> A2 + subgraph RESTConsumer["Web / Mobile / Dashboard"] + REST["REST-to-gRPC Gateway\n(apps/rest_gateway)\n(FastAPI - Port 8000)"] + end + + ESP32 -->|"MQTT Publish (JSON)"| Broker + Broker -->|"MQTT Subscribe"| Bridge + Bridge ==>|"gRPC ReportReading (HTTP/2)"| Coll + REST ==>|"gRPC Unary & Batch (HTTP/2)"| Coll ``` ### Industry-Grade Capabilities Implemented -- **Dual-App Separation**: Telemetry Collector (Server) and IoT Gateway (Client) operate as decoupled microservice applications. +- **Industrial IoT Protocol Hierarchy**: + - **MQTT**: Lightweight pub/sub for resource-constrained microcontrollers (ESP32) reading PZEM-004T sensors. + - **MQTT-to-gRPC Ingestion Bridge**: Seamlessly consumes MQTT telemetry topics and bridges them into gRPC. + - **FastAPI REST Gateway**: Modular HTTP backend for external web/mobile dashboards and REST API consumers. - **gRPC Interceptors**: - **Client-Side**: Injects distributed tracing headers (`x-request-id`) and `x-client-version`. - **Server-Side**: Performance metrics logging (RPC duration, peer IP, status code) and unhandled exception recovery translating errors safely into gRPC status codes. - **Official Health Checking (`grpc.health.v1`)**: Exposes standard gRPC health checks for Kubernetes liveness/readiness probes and load balancers. - **Connection Resilience**: Configured HTTP/2 keepalive pings (`grpc.keepalive_time_ms`), request timeouts (deadlines), and auto-reconnects. -- **REST-to-gRPC Gateway**: FastAPI application translating external HTTP/JSON REST requests into strongly-typed gRPC calls. -- **Container Orchestration**: Production `Dockerfile` and `docker-compose.yml` for multi-container deployment. +- **Container Orchestration**: Multi-container Docker Compose topology orchestrating Mosquitto MQTT, Collector, MQTT Bridge, and REST Gateway across segmented bridge networks. ### The four gRPC call types @@ -72,15 +74,22 @@ src/python_grpc/ device/ pzem_004t.py # PZEM004TDevice hardware physics simulator apps/ # autonomous deployable applications - collector/ # Cloud-tier: Telemetry Collector Server + collector/ # Cloud-tier: Telemetry Collector Server (gRPC only) app.py # server lifecycle, health check, graceful shutdown servicer.py # in-memory pub/sub telemetry broadcast servicer __main__.py # CLI entry point (python -m python_grpc.apps.collector) - device_agent/ # Edge-tier: IoT Hardware Agent - agent.py # resilient reporting, health checking, streaming - __main__.py # CLI entry point (python -m python_grpc.apps.device_agent) - rest_gateway/ # Consumer-tier: FastAPI REST-to-gRPC Gateway - app.py # FastAPI proxy endpoints (unary & batch) + mqtt_sensor_node/ # Edge-tier: Microcontroller (ESP32) MQTT Sensor Node + app.py # sensor reading & MQTT JSON publishing loop + __main__.py # CLI entry point (python -m python_grpc.apps.mqtt_sensor_node) + mqtt_bridge/ # Bridge-tier: MQTT-to-gRPC Telemetry Ingestion Bridge + app.py # MQTT subscriber forwarding to collector over gRPC + __main__.py # CLI entry point (python -m python_grpc.apps.mqtt_bridge) + rest_gateway/ # Consumer-tier: Modular FastAPI REST-to-gRPC Gateway + app.py # FastAPI app factory & lifespan + config.py # Gateway configuration settings + dependencies.py # Dependency injection & stub resolver + schemas.py # Pydantic models for validation + routers/ # Modular APIRouters (telemetry, health) __main__.py # CLI entry point (python -m python_grpc.apps.rest_gateway) scripts/ gen_proto.py # regenerate stubs from the proto @@ -91,8 +100,9 @@ tests/ test_cross_host.py # multi-client pub/sub broadcasting & disconnect resilience test_health_and_interceptors.py # gRPC health & interceptor integration tests test_gateway.py # FastAPI REST-to-gRPC gateway integration tests + test_mqtt_pipeline.py # MQTT sensor node + bridge + gRPC collector tests Dockerfile # multi-app container build -docker-compose.yml # multi-network orchestration (cloud-tier & edge-tier) +docker-compose.yml # multi-network orchestration with Mosquitto MQTT broker ``` ## Requirements @@ -116,24 +126,29 @@ poetry install poetry run python -m python_grpc.apps.collector --host 0.0.0.0 --port 50051 ``` -#### 2. Run the IoT Device Gateway Client (Terminal 2) -```bash -poetry run python -m python_grpc.apps.device_agent --target localhost:50051 --device-id PZEM-004T-0001 --count 5 -``` - -#### 3. Start the REST-to-gRPC Gateway (Terminal 3, optional) +#### 2. Start the REST-to-gRPC Gateway (Terminal 2) ```bash poetry run python -m python_grpc.apps.rest_gateway --host 0.0.0.0 --port 8000 --grpc-target localhost:50051 ``` Open your browser at `http://localhost:8000/docs` to test Swagger UI or send a cURL request: ```bash -curl -X POST "http://localhost:8000/api/v1/telemetry" \ +curl -X POST "http://localhost:8000/api/telemetry" \ -H "Content-Type: application/json" \ -d '{"device_id": "REST-01", "voltage": 230.2, "current": 2.1, "active_power": 483.4, "energy": 1.2, "frequency": 50.0, "power_factor": 0.99}' ``` +#### 3. Run the MQTT-to-gRPC Bridge & Simulated ESP32 Sensor Node (Terminal 3 & 4) +If you have an MQTT broker running (such as Mosquitto on port 1883): +```bash +# Start Bridge to forward MQTT messages into gRPC Collector +poetry run python -m python_grpc.apps.mqtt_bridge --mqtt-host localhost --mqtt-port 1883 --grpc-target localhost:50051 + +# Start Simulated ESP32 reading PZEM-004T and publishing over MQTT +poetry run python -m python_grpc.apps.mqtt_sensor_node --broker-host localhost --broker-port 1883 --device-id ESP32-PZEM-01 --count 5 +``` + ### Option B: Running with Docker Compose -Spin up the entire microservice topology: +Spin up the entire microservice topology (Mosquitto MQTT broker, Collector, MQTT Bridge, and REST Gateway): ```bash docker compose up --build ``` @@ -144,12 +159,30 @@ docker compose up --build poetry run pytest -v --cov=python_grpc ``` -Tests run 17 automated integration and unit tests covering: +Tests run 18 automated integration and unit tests covering: - Sensor physics and energy accumulation ([`tests/test_device.py`](tests/test_device.py)) - All 4 gRPC streaming patterns ([`tests/test_telemetry.py`](tests/test_telemetry.py)) - Cross-host multi-client pub/sub broadcasting & disconnect resilience ([`tests/test_cross_host.py`](tests/test_cross_host.py)) - Standard gRPC Health Checking (`grpc.health.v1`) & Interceptors ([`tests/test_health_and_interceptors.py`](tests/test_health_and_interceptors.py)) - FastAPI REST-to-gRPC unary and batch forwarding ([`tests/test_gateway.py`](tests/test_gateway.py)) +- Simulated ESP32 MQTT pub/sub ingestion into gRPC Collector ([`tests/test_mqtt_pipeline.py`](tests/test_mqtt_pipeline.py)) + +## Code Formatting & Linting + +Ruff is used for ultra-fast linting and code formatting: +```bash +# Check for lint violations +poetry run ruff check . + +# Automatically apply safe lint fixes +poetry run ruff check --fix . + +# Check formatting without modifying files +poetry run ruff format --check . + +# Automatically format the entire codebase +poetry run ruff format . +``` ## Type Checking diff --git a/docker-compose.yml b/docker-compose.yml index 56c193a..c9e1224 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -18,40 +18,74 @@ services: - collector.internal restart: unless-stopped + # ------------------------------------------------------------- - # Edge Device Tier (Simulates Remote IoT Edge Unit on Host 2) + # External Consumer Tier (Simulates REST Consumer on Host 3) # ------------------------------------------------------------- - device-agent: + rest-gateway: image: ${IMAGE_NAME:-python-grpc}:${IMAGE_TAG:-latest} build: context: . dockerfile: Dockerfile - command: python -m python_grpc.apps.device_agent --target collector.internal:50051 --device-id PZEM-EDGE-001 --count 5 + command: python -m python_grpc.apps.rest_gateway --host 0.0.0.0 --port 8000 --grpc-target collector.internal:50051 + ports: + - "8000:8000" depends_on: - collector-service environment: - - LOG_LEVEL=${LOG_LEVEL:-INFO} + - GRPC_TARGET=collector.internal:50051 + networks: + cloud-tier: + restart: unless-stopped + + # ------------------------------------------------------------- + # MQTT Broker (Industrial Standard IoT Messaging Tier) + # ------------------------------------------------------------- + mqtt-broker: + image: eclipse-mosquitto:2.0 + command: mosquitto -c /mosquitto-no-auth.conf + ports: + - "1883:1883" networks: edge-tier: + aliases: + - broker.internal cloud-tier: restart: unless-stopped # ------------------------------------------------------------- - # External Consumer Tier (Simulates REST Consumer on Host 3) + # IoT Microcontroller Sensor Node (Simulates ESP32 + PZEM-004T) # ------------------------------------------------------------- - rest-gateway: + mqtt-sensor-node: image: ${IMAGE_NAME:-python-grpc}:${IMAGE_TAG:-latest} build: context: . dockerfile: Dockerfile - command: python -m python_grpc.apps.rest_gateway --host 0.0.0.0 --port 8000 --grpc-target collector.internal:50051 - ports: - - "8000:8000" + command: python -m python_grpc.apps.mqtt_sensor_node --broker-host broker.internal --broker-port 1883 --device-id ESP32-PZEM-01 --count 5 + depends_on: + - mqtt-broker + environment: + - LOG_LEVEL=${LOG_LEVEL:-INFO} + networks: + edge-tier: + restart: unless-stopped + + # ------------------------------------------------------------- + # MQTT-to-gRPC Bridge Service (Translates MQTT pub/sub to gRPC) + # ------------------------------------------------------------- + mqtt-bridge: + image: ${IMAGE_NAME:-python-grpc}:${IMAGE_TAG:-latest} + build: + context: . + dockerfile: Dockerfile + command: python -m python_grpc.apps.mqtt_bridge --mqtt-host broker.internal --mqtt-port 1883 --grpc-target collector.internal:50051 depends_on: + - mqtt-broker - collector-service environment: - - GRPC_TARGET=collector.internal:50051 + - LOG_LEVEL=${LOG_LEVEL:-INFO} networks: + edge-tier: cloud-tier: restart: unless-stopped diff --git a/poetry.lock b/poetry.lock index c3f0322..359cf8f 100644 --- a/poetry.lock +++ b/poetry.lock @@ -484,6 +484,21 @@ files = [ {file = "packaging-26.3.tar.gz", hash = "sha256:94edc256424af38762eb31306eed28beb9f0efc50a8837492c9d6fd6004aed79"}, ] +[[package]] +name = "paho-mqtt" +version = "2.1.0" +description = "MQTT version 5.0/3.1.1 client class" +optional = false +python-versions = ">=3.7" +groups = ["main"] +files = [ + {file = "paho_mqtt-2.1.0-py3-none-any.whl", hash = "sha256:6db9ba9b34ed5bc6b6e3812718c7e06e2fd7444540df2455d2c51bd58808feee"}, + {file = "paho_mqtt-2.1.0.tar.gz", hash = "sha256:12d6e7511d4137555a3f6ea167ae846af2c7357b10bc6fa4f7c3968fc1723834"}, +] + +[package.extras] +proxy = ["pysocks"] + [[package]] name = "pluggy" version = "1.6.0" @@ -771,6 +786,34 @@ pytest = ">=7" [package.extras] testing = ["process-tests", "pytest-xdist", "virtualenv"] +[[package]] +name = "ruff" +version = "0.16.7" +description = "An extremely fast Python linter and code formatter, written in Rust." +optional = false +python-versions = ">=3.7" +groups = ["dev"] +files = [ + {file = "ruff-0.16.7-py3-none-linux_armv6l.whl", hash = "sha256:727307773e7c7f9181d3ed3a2484186e56c1fa1874255911c74585eb2c7c19f9"}, + {file = "ruff-0.16.7-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:9d61c258deabf58f34c67bd4bb4d939c7f2e6b5f0e59c1cdd1cf771b11cde929"}, + {file = "ruff-0.16.7-py3-none-macosx_11_0_arm64.whl", hash = "sha256:7ab81118df8945e0193d0240712aa4496573595b75185c3636ed825592a0f728"}, + {file = "ruff-0.16.7-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:4c196c968874fc8019da8e7163de7a1a370f111e2309b4b7dfea0fce950198d0"}, + {file = "ruff-0.16.7-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:ac8c3bd0a7e10ad31e6ce51e7a99f3cb772e69aecdd6b9ea7e99b362f62a62c0"}, + {file = "ruff-0.16.7-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:398d3988edde000b5c75dc1b3f584708da9bc990de069c18909142580fec1af9"}, + {file = "ruff-0.16.7-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:ce05b62b770a8217c4646a9c4139fca00efe8fe5d71f87df2b243ff20d4584d1"}, + {file = "ruff-0.16.7-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:af1b576fddb9d9ef2ececfb5fadcd6a624b25070ed85e3cfcfe449fc3ff6a7b9"}, + {file = "ruff-0.16.7-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:9ce7f8f22df67c93ed96c717f9128eadb797144ac2bad475cf536f31d6100c55"}, + {file = "ruff-0.16.7-py3-none-manylinux_2_31_riscv64.whl", hash = "sha256:06d0e93d04f392996435ebd600c153f65b47d73fbec2415aa99c5ee5756b3a5f"}, + {file = "ruff-0.16.7-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:142151a5e7b93c1b11111337142f89dd2fbfee92161225c99a97222f22e32656"}, + {file = "ruff-0.16.7-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:e6651f97a342d8b35d54d8991544ca22169b86dc54111cb604666940c431b750"}, + {file = "ruff-0.16.7-py3-none-musllinux_1_2_i686.whl", hash = "sha256:ef140c6eb935fa9a84c9c607dfb2cb1b85843c192e79265b0c54f35f557ea8e5"}, + {file = "ruff-0.16.7-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:53e39506a730fadeee0d998ed5946f30671f0db240c6c7c73bdabbe33604bb6f"}, + {file = "ruff-0.16.7-py3-none-win32.whl", hash = "sha256:2ea3470fcebcbc5df2fb0c6f3b90333fa9084c534e0111c038fa4a6ab9f1c4b7"}, + {file = "ruff-0.16.7-py3-none-win_amd64.whl", hash = "sha256:7ac26aca826e9e21d0f1cb25b54ac660760a9fdd094d3e4df9848232be98cfc6"}, + {file = "ruff-0.16.7-py3-none-win_arm64.whl", hash = "sha256:aab7f39e2c9df6c596216070f98eef1207b94f8516cca20c808826974971855b"}, + {file = "ruff-0.16.7.tar.gz", hash = "sha256:5f71d004ac1263b22fa39462ac5ae618a4b77d58981af2cc79bf79a29c12b1a6"}, +] + [[package]] name = "setuptools" version = "83.0.0" @@ -822,6 +865,22 @@ files = [ {file = "types_grpcio-1.83.0.20260730.tar.gz", hash = "sha256:e2b0532526a563dbb39aa6ee39fe0961488116ac1afec3065607532c93299b64"}, ] +[[package]] +name = "types-grpcio-health-checking" +version = "1.0.0.20260518" +description = "Typing stubs for grpcio-health-checking" +optional = false +python-versions = ">=3.10" +groups = ["dev"] +files = [ + {file = "types_grpcio_health_checking-1.0.0.20260518-py3-none-any.whl", hash = "sha256:e6b90d6cc7cf6509153be50384bed32f491b34f894d7ddb84e1c1bd0840ced8b"}, + {file = "types_grpcio_health_checking-1.0.0.20260518.tar.gz", hash = "sha256:05e402e439c3e1045b3e84af00ea5415ffda16497c89981e62c4d78c99aacb05"}, +] + +[package.dependencies] +types-grpcio = "*" +types-protobuf = "*" + [[package]] name = "types-protobuf" version = "7.34.1.20260518" @@ -882,5 +941,5 @@ standard = ["httptools (>=0.8.0)", "python-dotenv (>=0.13)", "pyyaml (>=5.1)", " [metadata] lock-version = "2.1" -python-versions = ">=3.13" -content-hash = "29d13ac8b94056ea12dca299d9d390458f6b70318bc11e6ac85de8098abed575" +python-versions = ">=3.13,<4.0" +content-hash = "c6b69fa3a2748d922bf011f93ae12d6f47e6aacfd3f7da309d5b75e92e81a5ed" diff --git a/pyproject.toml b/pyproject.toml index 74151eb..fd42641 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -6,7 +6,7 @@ authors = [ {name = "Farrel Augusta Dinata",email = "farrel.apeiron@gmail.com"} ] readme = "README.md" -requires-python = ">=3.13" +requires-python = ">=3.13,<4.0" dependencies = [ "fastapi (>=0.141.1,<0.142.0)", "grpcio (>=1.83.0,<2.0.0)", @@ -14,6 +14,7 @@ dependencies = [ "protobuf (>=7.35.1,<8.0.0)", "uvicorn (>=0.52.4,<0.53.0)", "grpcio-health-checking (>=1.83.1,<2.0.0)", + "paho-mqtt (>=2.1.0,<3.0.0)", ] [tool.poetry] @@ -36,15 +37,43 @@ build-backend = "poetry.core.masonry.api" [dependency-groups] dev = [ "pyrefly (>=1.2.0,<2.0.0)", + "ruff (>=0.14.0,<1.0.0)", "types-grpcio (>=1.83.0.20260730,<2.0.0.0)", "types-protobuf (>=7.34.1.20260518,<8.0.0.0)", "grpcio-health-checking (>=1.83.1,<2.0.0)", "httpx (>=0.28.1,<0.29.0)", "pytest (>=9.1.1,<10.0.0)", "pytest-asyncio (>=1.4.0,<2.0.0)", - "pytest-cov (>=7.1.0,<8.0.0)" + "pytest-cov (>=7.1.0,<8.0.0)", + "types-grpcio-health-checking (>=1.0.0.20260518,<2.0.0.0)" ] +[tool.ruff] +target-version = "py313" +line-length = 88 +extend-exclude = [ + "src/python_grpc/proto/*_pb2*.py", + "src/python_grpc/proto/*_pb2*.pyi", +] + +[tool.ruff.lint] +select = [ + "E", # pycodestyle errors + "W", # pycodestyle warnings + "F", # Pyflakes + "I", # isort + "B", # flake8-bugbear + "UP", # pyupgrade +] +ignore = [ + "E501", # line-too-long handled by ruff format +] + +[tool.ruff.format] +quote-style = "double" +indent-style = "space" + + [tool.pytest.ini_options] minversion = "8.0" testpaths = ["tests"] diff --git a/scripts/simulate_cross_host.py b/scripts/simulate_cross_host.py index 4a29346..0ed9316 100644 --- a/scripts/simulate_cross_host.py +++ b/scripts/simulate_cross_host.py @@ -3,7 +3,7 @@ Simulates two or more physically separate hosts communicating over gRPC: - Host 1 (Cloud Telemetry Collector Server): Runs gRPC server + health checks - Host 2 (Remote Subscriber / Dashboard Client): Subscribes to live readings stream - - Host 3 (Edge IoT Meter Gateway): Pushes periodic telemetry readings over gRPC + - Host 3 (Telemetry Ingestion Client / Bridge): Pushes periodic telemetry readings over gRPC Demonstrates that neither host shares memory or Python modules with the others; communication is 100% over the wire via HTTP/2 and Protobuf. @@ -13,18 +13,18 @@ import asyncio import logging -import time import grpc from python_grpc.apps.collector.servicer import DeviceTelemetryServicer -from python_grpc.apps.device_agent.agent import check_health, run_unary from python_grpc.core.common.config import DEFAULT_GRPC_CHANNEL_OPTIONS from python_grpc.core.common.interceptors import RequestIdClientInterceptor from python_grpc.core.device.pzem_004t import PZEM004TDevice from python_grpc.proto import pzem_004t_pb2, pzem_004t_pb2_grpc -logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] (%(name)s) %(message)s") +logging.basicConfig( + level=logging.INFO, format="%(asctime)s [%(levelname)s] (%(name)s) %(message)s" +) logger = logging.getLogger("cross_host_sim") @@ -56,8 +56,12 @@ async def host2_subscriber() -> None: interceptors=[RequestIdClientInterceptor(client_version="monitor-1.0")], ) as channel: stub = pzem_004t_pb2_grpc.DeviceTelemetryStub(channel) - host2_logger.info("Host 2 connected. Subscribing to live telemetry for 'PZEM-REMOTE-01'...") - call = stub.Subscribe(pzem_004t_pb2.SubscribeRequest(device_id="PZEM-REMOTE-01")) + host2_logger.info( + "Host 2 connected. Subscribing to live telemetry for 'PZEM-REMOTE-01'..." + ) + call = stub.Subscribe( + pzem_004t_pb2.SubscribeRequest(device_id="PZEM-REMOTE-01") + ) try: async for reading in call: host2_logger.info( @@ -74,8 +78,8 @@ async def host2_subscriber() -> None: except asyncio.CancelledError: pass - # 3. Spawn Host 3 (Edge IoT Meter Gateway Device) - host3_logger = logging.getLogger("Host3-IoTEdgeMeter") + # 3. Spawn Host 3 (Telemetry Ingestion Client / Bridge Service) + host3_logger = logging.getLogger("Host3-IngestionClient") async def host3_edge_device() -> None: await asyncio.sleep(0.2) # Give Host 2 time to establish subscription @@ -84,7 +88,7 @@ async def host3_edge_device() -> None: async with grpc.aio.insecure_channel( target, options=DEFAULT_GRPC_CHANNEL_OPTIONS, - interceptors=[RequestIdClientInterceptor(client_version="edge-iot-1.0")], + interceptors=[RequestIdClientInterceptor(client_version="ingestion-1.0")], ) as channel: stub = pzem_004t_pb2_grpc.DeviceTelemetryStub(channel) host3_logger.info("Host 3 connected to Cloud Server at %s", target) @@ -98,7 +102,9 @@ async def host3_edge_device() -> None: reading.current, ) ack = await stub.ReportReading(reading, timeout=5.0) - host3_logger.info("<- [ACKNOWLEDGED] Host 3 got server ack: %s", ack.message) + host3_logger.info( + "<- [ACKNOWLEDGED] Host 3 got server ack: %s", ack.message + ) await asyncio.sleep(0.3) # Run Host 2 and Host 3 concurrently communicating through Host 1 @@ -108,7 +114,9 @@ async def host3_edge_device() -> None: await asyncio.gather(t2, t3) print("\n" + "-" * 70) - print(f" Simulation complete: Host 2 received {len(received_readings)} live packets streamed from Host 3") + print( + f" Simulation complete: Host 2 received {len(received_readings)} live packets streamed from Host 3" + ) print("-" * 70 + "\n") await server.stop(grace=0) diff --git a/src/python_grpc/apps/collector/app.py b/src/python_grpc/apps/collector/app.py index 14b2430..076be8d 100644 --- a/src/python_grpc/apps/collector/app.py +++ b/src/python_grpc/apps/collector/app.py @@ -85,7 +85,9 @@ def _trigger_stop(*args: Any) -> None: def main() -> None: - parser = argparse.ArgumentParser(description="PZEM-004t gRPC Telemetry Collector Service") + parser = argparse.ArgumentParser( + description="PZEM-004t gRPC Telemetry Collector Service" + ) parser.add_argument( "--host", default=os.getenv("COLLECTOR_HOST", "0.0.0.0"), diff --git a/src/python_grpc/apps/collector/servicer.py b/src/python_grpc/apps/collector/servicer.py index f457155..953fd63 100644 --- a/src/python_grpc/apps/collector/servicer.py +++ b/src/python_grpc/apps/collector/servicer.py @@ -12,9 +12,9 @@ from __future__ import annotations import asyncio -from collections import defaultdict import logging import time +from collections import defaultdict from collections.abc import AsyncIterator from typing import override @@ -59,7 +59,9 @@ class DeviceTelemetryServicer(pzem_004t_pb2_grpc.DeviceTelemetryServicer): def __init__(self) -> None: self._readings: list[pzem_004t_pb2.ReadingReport] = [] # Pub/Sub registry: device_id -> set of active subscriber asyncio.Queues - self._subscribers: dict[str, set[asyncio.Queue[pzem_004t_pb2.ReadingReport]]] = defaultdict(set) + self._subscribers: dict[ + str, set[asyncio.Queue[pzem_004t_pb2.ReadingReport]] + ] = defaultdict(set) @property def readings(self) -> list[pzem_004t_pb2.ReadingReport]: @@ -73,13 +75,17 @@ def _broadcast(self, report: pzem_004t_pb2.ReadingReport) -> None: try: q.put_nowait(report) except asyncio.QueueFull: - logger.warning("Subscriber queue full, dropping reading for %s", report.device_id) + logger.warning( + "Subscriber queue full, dropping reading for %s", report.device_id + ) @override async def ReportReading( self, request: pzem_004t_pb2.ReadingReport, - context: grpc.aio.ServicerContext[pzem_004t_pb2.ReadingReport, pzem_004t_pb2.Ack], + context: grpc.aio.ServicerContext[ + pzem_004t_pb2.ReadingReport, pzem_004t_pb2.Ack + ], ) -> pzem_004t_pb2.Ack: if not _is_valid(request): return _ack(False, f"rejected invalid reading from {request.device_id}") @@ -92,13 +98,17 @@ async def ReportReading( request.current, request.active_power, ) - return _ack(True, f"stored reading #{len(self._readings)} from {request.device_id}") + return _ack( + True, f"stored reading #{len(self._readings)} from {request.device_id}" + ) @override async def ReportReadings( self, request_iterator: grpc.aio.AsyncIterable[pzem_004t_pb2.ReadingReport], - context: grpc.aio.ServicerContext[pzem_004t_pb2.ReadingReport, pzem_004t_pb2.BatchSummary], + context: grpc.aio.ServicerContext[ + pzem_004t_pb2.ReadingReport, pzem_004t_pb2.BatchSummary + ], ) -> pzem_004t_pb2.BatchSummary: received = 0 rejected = 0 @@ -128,7 +138,9 @@ async def ReportReadings( async def Subscribe( self, request: pzem_004t_pb2.SubscribeRequest, - context: grpc.aio.ServicerContext[pzem_004t_pb2.SubscribeRequest, pzem_004t_pb2.ReadingReport], + context: grpc.aio.ServicerContext[ + pzem_004t_pb2.SubscribeRequest, pzem_004t_pb2.ReadingReport + ], ) -> AsyncIterator[pzem_004t_pb2.ReadingReport]: target_device = request.device_id or "*" q: asyncio.Queue[pzem_004t_pb2.ReadingReport] = asyncio.Queue(maxsize=100) @@ -151,7 +163,9 @@ async def Subscribe( async def StreamTelemetry( self, request_iterator: grpc.aio.AsyncIterable[pzem_004t_pb2.ReadingReport], - context: grpc.aio.ServicerContext[pzem_004t_pb2.ReadingReport, pzem_004t_pb2.Ack], + context: grpc.aio.ServicerContext[ + pzem_004t_pb2.ReadingReport, pzem_004t_pb2.Ack + ], ) -> AsyncIterator[pzem_004t_pb2.Ack]: async for report in request_iterator: if _is_valid(report): diff --git a/src/python_grpc/apps/device_agent/__init__.py b/src/python_grpc/apps/device_agent/__init__.py deleted file mode 100644 index ad9d6dd..0000000 --- a/src/python_grpc/apps/device_agent/__init__.py +++ /dev/null @@ -1,25 +0,0 @@ -"""Device Agent application package.""" - -from python_grpc.apps.device_agent.agent import ( - DEFAULT_TARGET, - DEFAULT_TIMEOUT_S, - check_health, - run, - run_bidi, - run_client_streaming, - run_demo, - run_server_streaming, - run_unary, -) - -__all__ = [ - "DEFAULT_TARGET", - "DEFAULT_TIMEOUT_S", - "check_health", - "run", - "run_bidi", - "run_client_streaming", - "run_demo", - "run_server_streaming", - "run_unary", -] diff --git a/src/python_grpc/apps/device_agent/__main__.py b/src/python_grpc/apps/device_agent/__main__.py deleted file mode 100644 index 95b1241..0000000 --- a/src/python_grpc/apps/device_agent/__main__.py +++ /dev/null @@ -1,31 +0,0 @@ -"""CLI entry point for the PZEM-004t gRPC demo client.""" - -from __future__ import annotations - -import argparse - -from python_grpc.apps.device_agent.agent import DEFAULT_TARGET, run - - -def _parse_args() -> argparse.Namespace: - parser = argparse.ArgumentParser(description="PZEM-004t gRPC IoT Device Agent") - parser.add_argument("--target", default=DEFAULT_TARGET, help="server address, e.g. localhost:50051") - parser.add_argument("--device-id", default="PZEM-004T-0001") - parser.add_argument("--count", type=int, default=5, help="readings per streaming RPC") - parser.add_argument( - "--rpc", - default=["unary", "client-stream", "server-stream", "bidi"], - choices=["unary", "client-stream", "server-stream", "bidi"], - nargs="*", - help="RPC types to run (default: all)", - ) - return parser.parse_args() - - -def main() -> None: - args = _parse_args() - run(args.target, args.device_id, args.count, set(args.rpc)) - - -if __name__ == "__main__": - main() diff --git a/src/python_grpc/apps/device_agent/agent.py b/src/python_grpc/apps/device_agent/agent.py deleted file mode 100644 index 9232720..0000000 --- a/src/python_grpc/apps/device_agent/agent.py +++ /dev/null @@ -1,160 +0,0 @@ -"""Industry-grade resilient gRPC telemetry client. - -Features: - - HTTP/2 Keepalive & reconnection channel options - - Client interceptor injecting request-id and client version - - Client-side timeouts (deadlines) on unary and streams - - Comprehensive logging and health-check verification -""" - -# pyrefly: ignore-errors[missing-attribute] - -from __future__ import annotations - -import asyncio -import logging -from typing import AsyncIterable - -import grpc -from grpc_health.v1 import health_pb2, health_pb2_grpc - -from python_grpc.core.common.config import DEFAULT_GRPC_CHANNEL_OPTIONS -from python_grpc.core.common.interceptors import RequestIdClientInterceptor -from python_grpc.core.device.pzem_004t import PZEM004TDevice -from python_grpc.proto import pzem_004t_pb2, pzem_004t_pb2_grpc - -DEFAULT_TARGET = "localhost:50051" -DEFAULT_TIMEOUT_S = 10.0 - -logger = logging.getLogger("telemetry_device_client") - -DeviceTelemetryStub = pzem_004t_pb2_grpc.DeviceTelemetryStub - - -def _fmt(report: pzem_004t_pb2.ReadingReport) -> str: - return ( - f"{report.device_id} U={report.voltage:.1f}V I={report.current:.2f}A " - f"P={report.active_power:.1f}W E={report.energy:.6f}kWh " - f"f={report.frequency:.2f}Hz PF={report.power_factor:.3f}" - ) - - -async def check_health(channel: grpc.aio.Channel, service_name: str = "") -> bool: - """Query the remote gRPC server's standard health service.""" - health_stub = health_pb2_grpc.HealthStub(channel) - try: - response = await health_stub.Check( - health_pb2.HealthCheckRequest(service=service_name), - timeout=3.0, - ) - is_healthy = response.status == health_pb2.HealthCheckResponse.SERVING - logger.info( - "Health check service=%r status=%s", - service_name, - health_pb2.HealthCheckResponse.ServingStatus.Name(response.status), - ) - return is_healthy - except grpc.RpcError as exc: - logger.warning("Health check failed: code=%s details=%s", exc.code(), exc.details()) - return False - - -async def run_unary(stub: DeviceTelemetryStub, device: PZEM004TDevice, timeout: float = DEFAULT_TIMEOUT_S) -> None: - reading = device.read() - ack = await stub.ReportReading(reading, timeout=timeout) - logger.info("[unary] ReportReading -> success=%s msg=%r", ack.success, ack.message) - - -async def run_client_streaming( - stub: DeviceTelemetryStub, - device: PZEM004TDevice, - count: int, - timeout: float = DEFAULT_TIMEOUT_S, -) -> None: - async def readings() -> AsyncIterable[pzem_004t_pb2.ReadingReport]: - for _ in range(count): - yield device.read() - - summary = await stub.ReportReadings(readings(), timeout=timeout) - logger.info( - "[client] ReportReadings -> received=%d rejected=%d avg_power=%.1fW", - summary.received, - summary.rejected, - summary.avg_active_power, - ) - - -async def run_server_streaming( - stub: DeviceTelemetryStub, - device: PZEM004TDevice, - count: int, -) -> None: - request = pzem_004t_pb2.SubscribeRequest(device_id=device.device_id) - logger.info("[server] Subscribe -> live readings for %s:", device.device_id) - call = stub.Subscribe(request) - i = 0 - try: - async for report in call: - logger.info(" %s", _fmt(report)) - i += 1 - if i >= count: - call.cancel() - break - except asyncio.CancelledError: - pass - - -async def run_bidi( - stub: DeviceTelemetryStub, - device: PZEM004TDevice, - count: int, -) -> None: - async def readings() -> AsyncIterable[pzem_004t_pb2.ReadingReport]: - for _ in range(count): - yield device.read() - - logger.info("[bidi] StreamTelemetry -> %d readings sent, acks:", count) - async for ack in stub.StreamTelemetry(readings()): - logger.info(" success=%s msg=%r", ack.success, ack.message) - - -async def run_demo( - target: str, - device_id: str, - count: int, - rpcs: set[str], - timeout: float = DEFAULT_TIMEOUT_S, -) -> None: - device = PZEM004TDevice(device_id=device_id, read_interval_s=1.0) - interceptors = [RequestIdClientInterceptor(client_version="2.0.0")] - - async with grpc.aio.insecure_channel( - target, - options=DEFAULT_GRPC_CHANNEL_OPTIONS, - interceptors=interceptors, - ) as channel: - # Verify server health first - logger.info("Connecting to collector at %s (Device: %s)", target, device_id) - is_healthy = await check_health(channel) - if not is_healthy: - logger.warning("Target %s is not currently reporting SERVING status.", target) - - stub = pzem_004t_pb2_grpc.DeviceTelemetryStub(channel) - - if "unary" in rpcs: - await run_unary(stub, device, timeout=timeout) - if "client-stream" in rpcs: - await run_client_streaming(stub, device, count, timeout=timeout) - if "server-stream" in rpcs: - await run_server_streaming(stub, device, count) - if "bidi" in rpcs: - await run_bidi(stub, device, count) - logger.info("Demo complete.") - - -def run(target: str, device_id: str, count: int, rpcs: set[str]) -> None: - logging.basicConfig( - level=logging.INFO, - format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", - ) - asyncio.run(run_demo(target, device_id, count, rpcs)) diff --git a/src/python_grpc/apps/mqtt_bridge/__init__.py b/src/python_grpc/apps/mqtt_bridge/__init__.py new file mode 100644 index 0000000..c3572c3 --- /dev/null +++ b/src/python_grpc/apps/mqtt_bridge/__init__.py @@ -0,0 +1,5 @@ +"""MQTT-to-gRPC Ingestion Bridge package.""" + +from python_grpc.apps.mqtt_bridge.app import bridge_mqtt_to_grpc + +__all__ = ["bridge_mqtt_to_grpc"] diff --git a/src/python_grpc/apps/mqtt_bridge/__main__.py b/src/python_grpc/apps/mqtt_bridge/__main__.py new file mode 100644 index 0000000..354e42c --- /dev/null +++ b/src/python_grpc/apps/mqtt_bridge/__main__.py @@ -0,0 +1,6 @@ +"""CLI entry point for running the MQTT-to-gRPC Ingestion Bridge.""" + +from python_grpc.apps.mqtt_bridge.app import main + +if __name__ == "__main__": + main() diff --git a/src/python_grpc/apps/mqtt_bridge/app.py b/src/python_grpc/apps/mqtt_bridge/app.py new file mode 100644 index 0000000..5bb609f --- /dev/null +++ b/src/python_grpc/apps/mqtt_bridge/app.py @@ -0,0 +1,181 @@ +"""MQTT-to-gRPC Bridge service. + +Subscribes to sensor telemetry topics on an MQTT broker using paho-mqtt +and bridges them into the central Telemetry Collector via high-performance gRPC. +""" + +from __future__ import annotations + +import argparse +import asyncio +import json +import logging +import os + +import grpc +import paho.mqtt.client as mqtt + +from python_grpc.core.common.config import DEFAULT_GRPC_CHANNEL_OPTIONS +from python_grpc.core.common.interceptors import RequestIdClientInterceptor +from python_grpc.proto import pzem_004t_pb2, pzem_004t_pb2_grpc + +logger = logging.getLogger("mqtt_grpc_bridge") + + +async def bridge_mqtt_to_grpc( + mqtt_host: str, + mqtt_port: int, + grpc_target: str, + topic: str = "devices/+/telemetry", + stop_after_count: int = 0, +) -> None: + """Subscribe to MQTT topic pattern via paho-mqtt and forward payloads into gRPC.""" + logger.info( + "Starting MQTT-to-gRPC Bridge (paho-mqtt): MQTT %s:%d (topic=%r) -> gRPC %s", + mqtt_host, + mqtt_port, + topic, + grpc_target, + ) + + channel = grpc.aio.insecure_channel( + grpc_target, + options=DEFAULT_GRPC_CHANNEL_OPTIONS, + interceptors=[RequestIdClientInterceptor(client_version="mqtt-bridge-1.0")], + ) + stub = pzem_004t_pb2_grpc.DeviceTelemetryStub(channel) + + loop = asyncio.get_running_loop() + msg_queue: asyncio.Queue[pzem_004t_pb2.ReadingReport] = asyncio.Queue() + + client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2) + + def on_connect( + client: mqtt.Client, + userdata: object, + flags: mqtt.ConnectFlags, + rc: mqtt.ReasonCode, + properties: mqtt.Properties | None = None, + ) -> None: + if rc.is_failure: + logger.error("Failed to connect to MQTT broker: %s", rc) + return + logger.info("Connected to MQTT broker. Subscribing to %s...", topic) + client.subscribe(topic) + + def on_message( + client: mqtt.Client, + userdata: object, + message: mqtt.MQTTMessage, + ) -> None: + try: + payload = json.loads(message.payload.decode("utf-8")) + reading = pzem_004t_pb2.ReadingReport( + device_id=payload.get("device_id", "UNKNOWN"), + device_type=payload.get("device_type", "PZEM-004T"), + timestamp_unix_ms=payload.get("timestamp_unix_ms", 0), + voltage=float(payload.get("voltage", 0.0)), + current=float(payload.get("current", 0.0)), + active_power=float(payload.get("active_power", 0.0)), + energy=float(payload.get("energy", 0.0)), + frequency=float(payload.get("frequency", 50.0)), + power_factor=float(payload.get("power_factor", 1.0)), + ) + loop.call_soon_threadsafe(msg_queue.put_nowait, reading) + except (json.JSONDecodeError, KeyError, ValueError) as err: + logger.warning("Failed to parse MQTT message payload: %s", err) + + client.on_connect = on_connect + client.on_message = on_message + + client.connect(mqtt_host, mqtt_port, keepalive=60) + client.loop_start() + + forwarded = 0 + try: + while True: + reading = await msg_queue.get() + try: + ack = await stub.ReportReading(reading, timeout=5.0) + forwarded += 1 + logger.info( + "[%d] Forwarded MQTT msg from %s to gRPC: %s", + forwarded, + reading.device_id, + ack.message, + ) + except grpc.RpcError as rpc_err: + logger.error("gRPC forward failed: %s", rpc_err.details()) + + if 0 < stop_after_count <= forwarded: + logger.info( + "Reached stop count %d, exiting bridge loop.", stop_after_count + ) + break + finally: + client.loop_stop() + client.disconnect() + await channel.close() + logger.info("Bridge gRPC channel and MQTT client closed.") + + +def main() -> None: + parser = argparse.ArgumentParser( + description="MQTT-to-gRPC Telemetry Ingestion Bridge (paho-mqtt)" + ) + parser.add_argument( + "--mqtt-host", + default=os.getenv("MQTT_BROKER_HOST", "localhost"), + help="MQTT broker host", + ) + parser.add_argument( + "--mqtt-port", + type=int, + default=int(os.getenv("MQTT_BROKER_PORT", "1883")), + help="MQTT broker port", + ) + parser.add_argument( + "--grpc-target", + default=os.getenv("GRPC_TARGET", "localhost:50051"), + help="Target gRPC Telemetry Collector address", + ) + parser.add_argument( + "--topic", + default=os.getenv("MQTT_TOPIC", "devices/+/telemetry"), + help="MQTT topic filter to subscribe to", + ) + parser.add_argument( + "--count", + type=int, + default=int(os.getenv("COUNT", "0")), + help="Stop after forwarding N messages (0 for continuous)", + ) + parser.add_argument( + "--log-level", + default=os.getenv("LOG_LEVEL", "INFO"), + choices=["DEBUG", "INFO", "WARNING", "ERROR"], + help="Logging level", + ) + args = parser.parse_args() + + logging.basicConfig( + level=getattr(logging, args.log_level), + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", + ) + + try: + asyncio.run( + bridge_mqtt_to_grpc( + mqtt_host=args.mqtt_host, + mqtt_port=args.mqtt_port, + grpc_target=args.grpc_target, + topic=args.topic, + stop_after_count=args.count, + ) + ) + except KeyboardInterrupt: + logger.info("Bridge stopped by user.") + + +if __name__ == "__main__": + main() diff --git a/src/python_grpc/apps/mqtt_sensor_node/__init__.py b/src/python_grpc/apps/mqtt_sensor_node/__init__.py new file mode 100644 index 0000000..9bb1898 --- /dev/null +++ b/src/python_grpc/apps/mqtt_sensor_node/__init__.py @@ -0,0 +1,5 @@ +"""Simulated ESP32 MQTT Sensor Node package.""" + +from python_grpc.apps.mqtt_sensor_node.app import run_sensor_node + +__all__ = ["run_sensor_node"] diff --git a/src/python_grpc/apps/mqtt_sensor_node/__main__.py b/src/python_grpc/apps/mqtt_sensor_node/__main__.py new file mode 100644 index 0000000..1e315c8 --- /dev/null +++ b/src/python_grpc/apps/mqtt_sensor_node/__main__.py @@ -0,0 +1,6 @@ +"""CLI runner for simulated ESP32 MQTT Sensor Node.""" + +from python_grpc.apps.mqtt_sensor_node.app import main + +if __name__ == "__main__": + main() diff --git a/src/python_grpc/apps/mqtt_sensor_node/app.py b/src/python_grpc/apps/mqtt_sensor_node/app.py new file mode 100644 index 0000000..2019220 --- /dev/null +++ b/src/python_grpc/apps/mqtt_sensor_node/app.py @@ -0,0 +1,152 @@ +"""Simulated Microcontroller (ESP32) reading PZEM-004T and publishing over MQTT.""" + +from __future__ import annotations + +import argparse +import json +import logging +import os +import time + +import paho.mqtt.client as mqtt + +from python_grpc.core.device.pzem_004t import PZEM004TDevice + +logger = logging.getLogger("mqtt_sensor_node") + + +def run_sensor_node( + broker_host: str, + broker_port: int, + device_id: str, + topic_prefix: str = "devices", + interval_s: float = 1.0, + count: int = 0, +) -> None: + """Read sensor data from PZEM-004T and publish as JSON over MQTT using standard paho-mqtt.""" + sensor = PZEM004TDevice(device_id=device_id, read_interval_s=interval_s) + topic = f"{topic_prefix}/{device_id}/telemetry" + + logger.info( + "Connecting to MQTT broker %s:%d to publish on topic %s", + broker_host, + broker_port, + topic, + ) + + client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2) + + def on_connect( + client: mqtt.Client, + userdata: object, + flags: mqtt.ConnectFlags, + rc: mqtt.ReasonCode, + properties: mqtt.Properties | None = None, + ) -> None: + if rc.is_failure: + logger.error("Failed to connect to MQTT broker: %s", rc) + else: + logger.info("Connected to MQTT broker successfully.") + + client.on_connect = on_connect + + try: + client.connect(broker_host, broker_port, keepalive=60) + client.loop_start() + + published = 0 + while True: + reading = sensor.read() + payload = { + "device_id": reading.device_id, + "device_type": reading.device_type, + "timestamp_unix_ms": reading.timestamp_unix_ms, + "voltage": reading.voltage, + "current": reading.current, + "active_power": reading.active_power, + "energy": reading.energy, + "frequency": reading.frequency, + "power_factor": reading.power_factor, + } + raw_payload = json.dumps(payload) + client.publish(topic, payload=raw_payload, qos=1) + published += 1 + logger.info( + "[%d] Published to %s: V=%.1fV I=%.2fA P=%.1fW", + published, + topic, + reading.voltage, + reading.current, + reading.active_power, + ) + + if 0 < count <= published: + logger.info("Published requested count (%d), stopping.", count) + break + + time.sleep(interval_s) + finally: + client.loop_stop() + client.disconnect() + logger.info("Disconnected from MQTT broker.") + + +def main() -> None: + parser = argparse.ArgumentParser( + description="Simulated ESP32 + PZEM-004T MQTT Sensor Node (paho-mqtt)" + ) + parser.add_argument( + "--broker-host", + default=os.getenv("MQTT_BROKER_HOST", "localhost"), + help="MQTT broker host", + ) + parser.add_argument( + "--broker-port", + type=int, + default=int(os.getenv("MQTT_BROKER_PORT", "1883")), + help="MQTT broker port", + ) + parser.add_argument( + "--device-id", + default=os.getenv("DEVICE_ID", "ESP32-PZEM-01"), + help="Device ID identifier", + ) + parser.add_argument( + "--interval", + type=float, + default=float(os.getenv("READ_INTERVAL_S", "1.0")), + help="Publishing interval in seconds", + ) + parser.add_argument( + "--count", + type=int, + default=int(os.getenv("COUNT", "0")), + help="Number of readings to publish (0 for continuous)", + ) + parser.add_argument( + "--log-level", + default=os.getenv("LOG_LEVEL", "INFO"), + choices=["DEBUG", "INFO", "WARNING", "ERROR"], + help="Logging level", + ) + args = parser.parse_args() + + logging.basicConfig( + level=getattr(logging, args.log_level), + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", + ) + + try: + run_sensor_node( + broker_host=args.broker_host, + broker_port=args.broker_port, + device_id=args.device_id, + interval_s=args.interval, + count=args.count, + ) + except KeyboardInterrupt: + logger.info("Sensor node interrupted by user.") + + +if __name__ == "__main__": + main() diff --git a/src/python_grpc/apps/rest_gateway/__init__.py b/src/python_grpc/apps/rest_gateway/__init__.py index 1f665fc..92e0bb3 100644 --- a/src/python_grpc/apps/rest_gateway/__init__.py +++ b/src/python_grpc/apps/rest_gateway/__init__.py @@ -1,5 +1,6 @@ """REST-to-gRPC Gateway application package.""" -from python_grpc.apps.rest_gateway.app import app +from python_grpc.apps.rest_gateway.app import app, create_app +from python_grpc.apps.rest_gateway.dependencies import state -__all__ = ["app"] +__all__ = ["app", "create_app", "state"] diff --git a/src/python_grpc/apps/rest_gateway/__main__.py b/src/python_grpc/apps/rest_gateway/__main__.py index edf8cf0..9144c75 100644 --- a/src/python_grpc/apps/rest_gateway/__main__.py +++ b/src/python_grpc/apps/rest_gateway/__main__.py @@ -4,19 +4,40 @@ import argparse import os + import uvicorn def main() -> None: - parser = argparse.ArgumentParser(description="PZEM-004t REST-to-gRPC Gateway Service") - parser.add_argument("--host", default=os.getenv("GATEWAY_HOST", "0.0.0.0"), help="Host to bind") - parser.add_argument("--port", type=int, default=int(os.getenv("GATEWAY_PORT", "8000")), help="Port to bind") - parser.add_argument("--grpc-target", default=os.getenv("GRPC_TARGET", "localhost:50051"), help="Target gRPC server") + parser = argparse.ArgumentParser( + description="PZEM-004t REST-to-gRPC Gateway Service" + ) + parser.add_argument( + "--host", default=os.getenv("GATEWAY_HOST", "0.0.0.0"), help="Host to bind" + ) + parser.add_argument( + "--port", + type=int, + default=int(os.getenv("GATEWAY_PORT", "8000")), + help="Port to bind", + ) + parser.add_argument( + "--grpc-target", + default=os.getenv("GRPC_TARGET", "localhost:50051"), + help="Target gRPC server", + ) args = parser.parse_args() os.environ["GRPC_TARGET"] = args.grpc_target - print(f"Starting REST-to-gRPC Gateway on http://{args.host}:{args.port} -> gRPC {args.grpc_target}") - uvicorn.run("python_grpc.apps.rest_gateway.app:app", host=args.host, port=args.port, reload=False) + print( + f"Starting REST-to-gRPC Gateway on http://{args.host}:{args.port} -> gRPC {args.grpc_target}" + ) + uvicorn.run( + "python_grpc.apps.rest_gateway.app:app", + host=args.host, + port=args.port, + reload=False, + ) if __name__ == "__main__": diff --git a/src/python_grpc/apps/rest_gateway/app.py b/src/python_grpc/apps/rest_gateway/app.py index 3b2ba8d..f3588b8 100644 --- a/src/python_grpc/apps/rest_gateway/app.py +++ b/src/python_grpc/apps/rest_gateway/app.py @@ -2,56 +2,22 @@ from __future__ import annotations -import os -from collections.abc import AsyncIterator +from collections.abc import AsyncGenerator from contextlib import asynccontextmanager -from typing import Any -from fastapi import FastAPI, HTTPException, Query import grpc -from pydantic import BaseModel, Field +from fastapi import FastAPI +from python_grpc.apps.rest_gateway.config import GRPC_TARGET +from python_grpc.apps.rest_gateway.dependencies import GatewayState, state +from python_grpc.apps.rest_gateway.routers import health_router, telemetry_router from python_grpc.core.common.config import DEFAULT_GRPC_CHANNEL_OPTIONS from python_grpc.core.common.interceptors import RequestIdClientInterceptor -from python_grpc.proto import pzem_004t_pb2, pzem_004t_pb2_grpc - -GRPC_TARGET = os.getenv("GRPC_TARGET", "localhost:50051") - - -class ReadingPayload(BaseModel): - device_id: str = Field(..., examples=["PZEM-004T-REST-01"]) - device_type: str = Field("PZEM-004T") - timestamp_unix_ms: int | None = None - voltage: float = Field(..., ge=0.0, le=1000.0) - current: float = Field(..., ge=0.0, le=1000.0) - active_power: float = Field(..., ge=0.0) - energy: float = Field(..., ge=0.0) - frequency: float = Field(..., ge=45.0, le=65.0) - power_factor: float = Field(..., ge=0.0, le=1.0) - - -class AckResponse(BaseModel): - success: bool - message: str - received_at_unix_ms: int - - -class BatchSummaryResponse(BaseModel): - received: int - rejected: int - avg_active_power: float - - -class GatewayState: - channel: grpc.aio.Channel | None = None - stub: pzem_004t_pb2_grpc.DeviceTelemetryStub | None = None - - -state = GatewayState() +from python_grpc.proto import pzem_004t_pb2_grpc @asynccontextmanager -async def lifespan(app: FastAPI) -> AsyncIterator[None]: +async def lifespan(app: FastAPI) -> AsyncGenerator[None]: # Initialize gRPC channel on startup state.channel = grpc.aio.insecure_channel( GRPC_TARGET, @@ -65,73 +31,19 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: await state.channel.close() -app = FastAPI( - title="PZEM-004t REST-to-gRPC Gateway", - description="Demonstrates REST API consumers communicating with internal gRPC microservices.", - version="1.0.0", - lifespan=lifespan, -) - - -@app.post("/api/v1/telemetry", response_model=AckResponse) -async def report_reading(payload: ReadingPayload) -> AckResponse: - """Proxy HTTP POST to gRPC ReportReading (unary).""" - if state.stub is None: - raise HTTPException(status_code=503, detail="gRPC gateway channel unavailable") - - request = pzem_004t_pb2.ReadingReport( - device_id=payload.device_id, - device_type=payload.device_type, - timestamp_unix_ms=payload.timestamp_unix_ms or 0, - voltage=payload.voltage, - current=payload.current, - active_power=payload.active_power, - energy=payload.energy, - frequency=payload.frequency, - power_factor=payload.power_factor, +def create_app() -> FastAPI: + """Factory creating and configuring the FastAPI Gateway application.""" + application = FastAPI( + title="PZEM-004t REST-to-gRPC Gateway", + description="Demonstrates REST API consumers communicating with internal gRPC microservices.", + version="1.0.0", + lifespan=lifespan, ) - try: - ack = await state.stub.ReportReading(request, timeout=5.0) - return AckResponse( - success=ack.success, - message=ack.message, - received_at_unix_ms=ack.received_at_unix_ms, - ) - except grpc.RpcError as exc: - raise HTTPException(status_code=502, detail=f"gRPC call failed: {exc.details()}") from exc - - -@app.post("/api/v1/telemetry/batch", response_model=BatchSummaryResponse) -async def report_batch(readings: list[ReadingPayload]) -> BatchSummaryResponse: - """Proxy HTTP POST list to gRPC ReportReadings (client-streaming).""" - if state.stub is None: - raise HTTPException(status_code=503, detail="gRPC gateway channel unavailable") - - async def _generator(): - for r in readings: - yield pzem_004t_pb2.ReadingReport( - device_id=r.device_id, - device_type=r.device_type, - timestamp_unix_ms=r.timestamp_unix_ms or 0, - voltage=r.voltage, - current=r.current, - active_power=r.active_power, - energy=r.energy, - frequency=r.frequency, - power_factor=r.power_factor, - ) + application.include_router(health_router) + application.include_router(telemetry_router) + return application - try: - summary = await state.stub.ReportReadings(_generator(), timeout=10.0) - return BatchSummaryResponse( - received=summary.received, - rejected=summary.rejected, - avg_active_power=summary.avg_active_power, - ) - except grpc.RpcError as exc: - raise HTTPException(status_code=502, detail=f"gRPC call failed: {exc.details()}") from exc +app = create_app() -@app.get("/health") -async def health_check() -> dict[str, str]: - return {"status": "ok", "grpc_target": GRPC_TARGET} +__all__ = ["GatewayState", "app", "create_app", "lifespan", "state"] diff --git a/src/python_grpc/apps/rest_gateway/config.py b/src/python_grpc/apps/rest_gateway/config.py new file mode 100644 index 0000000..ff7a598 --- /dev/null +++ b/src/python_grpc/apps/rest_gateway/config.py @@ -0,0 +1,9 @@ +"""Configuration settings for the REST-to-gRPC Gateway service.""" + +from __future__ import annotations + +import os + +GRPC_TARGET: str = os.getenv("GRPC_TARGET", "localhost:50051") +GATEWAY_HOST: str = os.getenv("GATEWAY_HOST", "0.0.0.0") +GATEWAY_PORT: int = int(os.getenv("GATEWAY_PORT", "8000")) diff --git a/src/python_grpc/apps/rest_gateway/dependencies.py b/src/python_grpc/apps/rest_gateway/dependencies.py new file mode 100644 index 0000000..de6fd2d --- /dev/null +++ b/src/python_grpc/apps/rest_gateway/dependencies.py @@ -0,0 +1,28 @@ +"""Dependency injection and gRPC connection management for the REST Gateway.""" + +from __future__ import annotations + +import grpc +from fastapi import HTTPException + +from python_grpc.proto import pzem_004t_pb2_grpc + + +class GatewayState: + """Manages the lifecycle of gRPC channel and client stubs.""" + + channel: grpc.aio.Channel | None = None + stub: pzem_004t_pb2_grpc.DeviceTelemetryStub | None = None + + +state = GatewayState() + + +def get_telemetry_stub() -> pzem_004t_pb2_grpc.DeviceTelemetryStub: + """FastAPI dependency yielding the active gRPC telemetry client stub.""" + if state.stub is None: + raise HTTPException( + status_code=503, + detail="gRPC gateway channel unavailable", + ) + return state.stub diff --git a/src/python_grpc/apps/rest_gateway/routers/__init__.py b/src/python_grpc/apps/rest_gateway/routers/__init__.py new file mode 100644 index 0000000..2d49518 --- /dev/null +++ b/src/python_grpc/apps/rest_gateway/routers/__init__.py @@ -0,0 +1,8 @@ +"""Routers package for the REST Gateway.""" + +from __future__ import annotations + +from python_grpc.apps.rest_gateway.routers.health import router as health_router +from python_grpc.apps.rest_gateway.routers.telemetry import router as telemetry_router + +__all__ = ["health_router", "telemetry_router"] diff --git a/src/python_grpc/apps/rest_gateway/routers/health.py b/src/python_grpc/apps/rest_gateway/routers/health.py new file mode 100644 index 0000000..438fbd6 --- /dev/null +++ b/src/python_grpc/apps/rest_gateway/routers/health.py @@ -0,0 +1,15 @@ +"""Health check endpoints for the REST Gateway.""" + +from __future__ import annotations + +from fastapi import APIRouter + +from python_grpc.apps.rest_gateway.config import GRPC_TARGET + +router = APIRouter(tags=["Health"]) + + +@router.get("/health") +async def health_check() -> dict[str, str]: + """Return health status and configured gRPC target.""" + return {"status": "ok", "grpc_target": GRPC_TARGET} diff --git a/src/python_grpc/apps/rest_gateway/routers/telemetry.py b/src/python_grpc/apps/rest_gateway/routers/telemetry.py new file mode 100644 index 0000000..68d63ae --- /dev/null +++ b/src/python_grpc/apps/rest_gateway/routers/telemetry.py @@ -0,0 +1,87 @@ +"""Telemetry REST endpoints forwarding payloads to internal gRPC collector.""" + +from __future__ import annotations + +from collections.abc import AsyncGenerator +from typing import Annotated + +import grpc +from fastapi import APIRouter, Depends, HTTPException + +from python_grpc.apps.rest_gateway.dependencies import get_telemetry_stub +from python_grpc.apps.rest_gateway.schemas import ( + AckResponse, + BatchSummaryResponse, + ReadingPayload, +) +from python_grpc.proto import pzem_004t_pb2, pzem_004t_pb2_grpc + +TelemetryStubDep = Annotated[ + pzem_004t_pb2_grpc.DeviceTelemetryStub, Depends(get_telemetry_stub) +] + +router = APIRouter(prefix="/api/telemetry", tags=["Telemetry"]) + + +@router.post("", response_model=AckResponse) +async def report_reading( + payload: ReadingPayload, + stub: TelemetryStubDep, +) -> AckResponse: + """Proxy HTTP POST to gRPC ReportReading (unary).""" + request = pzem_004t_pb2.ReadingReport( + device_id=payload.device_id, + device_type=payload.device_type, + timestamp_unix_ms=payload.timestamp_unix_ms or 0, + voltage=payload.voltage, + current=payload.current, + active_power=payload.active_power, + energy=payload.energy, + frequency=payload.frequency, + power_factor=payload.power_factor, + ) + try: + ack = await stub.ReportReading(request, timeout=5.0) + return AckResponse( + success=ack.success, + message=ack.message, + received_at_unix_ms=ack.received_at_unix_ms, + ) + except grpc.RpcError as exc: + raise HTTPException( + status_code=502, detail=f"gRPC call failed: {exc.details()}" + ) from exc + + +@router.post("/batch", response_model=BatchSummaryResponse) +async def report_batch( + readings: list[ReadingPayload], + stub: TelemetryStubDep, +) -> BatchSummaryResponse: + """Proxy HTTP POST list to gRPC ReportReadings (client-streaming).""" + + async def _generator() -> AsyncGenerator[pzem_004t_pb2.ReadingReport]: + for r in readings: + yield pzem_004t_pb2.ReadingReport( + device_id=r.device_id, + device_type=r.device_type, + timestamp_unix_ms=r.timestamp_unix_ms or 0, + voltage=r.voltage, + current=r.current, + active_power=r.active_power, + energy=r.energy, + frequency=r.frequency, + power_factor=r.power_factor, + ) + + try: + summary = await stub.ReportReadings(_generator(), timeout=10.0) + return BatchSummaryResponse( + received=summary.received, + rejected=summary.rejected, + avg_active_power=summary.avg_active_power, + ) + except grpc.RpcError as exc: + raise HTTPException( + status_code=502, detail=f"gRPC call failed: {exc.details()}" + ) from exc diff --git a/src/python_grpc/apps/rest_gateway/schemas.py b/src/python_grpc/apps/rest_gateway/schemas.py new file mode 100644 index 0000000..14d3f1b --- /dev/null +++ b/src/python_grpc/apps/rest_gateway/schemas.py @@ -0,0 +1,29 @@ +"""Pydantic data validation schemas for the REST-to-gRPC Gateway.""" + +from __future__ import annotations + +from pydantic import BaseModel, Field + + +class ReadingPayload(BaseModel): + device_id: str = Field(..., examples=["PZEM-004T-REST-01"]) + device_type: str = Field("PZEM-004T") + timestamp_unix_ms: int | None = None + voltage: float = Field(..., ge=0.0, le=1000.0) + current: float = Field(..., ge=0.0, le=1000.0) + active_power: float = Field(..., ge=0.0) + energy: float = Field(..., ge=0.0) + frequency: float = Field(..., ge=45.0, le=65.0) + power_factor: float = Field(..., ge=0.0, le=1.0) + + +class AckResponse(BaseModel): + success: bool + message: str + received_at_unix_ms: int + + +class BatchSummaryResponse(BaseModel): + received: int + rejected: int + avg_active_power: float diff --git a/src/python_grpc/core/common/config.py b/src/python_grpc/core/common/config.py index 5820cb0..70f6197 100644 --- a/src/python_grpc/core/common/config.py +++ b/src/python_grpc/core/common/config.py @@ -6,10 +6,13 @@ # Production HTTP/2 channel options for keepalive and resiliency DEFAULT_GRPC_CHANNEL_OPTIONS: list[tuple[str, Any]] = [ - ("grpc.keepalive_time_ms", 30000), # Send keepalive ping every 30s - ("grpc.keepalive_timeout_ms", 10000), # Keepalive ping timeout 10s - ("grpc.keepalive_permit_without_calls", 1), # Allow pings even when no active RPC calls - ("grpc.http2.max_pings_without_data", 0), # Unlimited pings without data + ("grpc.keepalive_time_ms", 30000), # Send keepalive ping every 30s + ("grpc.keepalive_timeout_ms", 10000), # Keepalive ping timeout 10s + ( + "grpc.keepalive_permit_without_calls", + 1, + ), # Allow pings even when no active RPC calls + ("grpc.http2.max_pings_without_data", 0), # Unlimited pings without data ("grpc.http2.min_time_between_pings_ms", 10000), ("grpc.http2.min_ping_interval_without_data_ms", 5000), ] diff --git a/src/python_grpc/core/common/interceptors.py b/src/python_grpc/core/common/interceptors.py index 57300b4..aa2c1e4 100644 --- a/src/python_grpc/core/common/interceptors.py +++ b/src/python_grpc/core/common/interceptors.py @@ -24,13 +24,17 @@ CLIENT_VERSION_HEADER = "x-client-version" -class RequestIdClientInterceptor(UnaryUnaryClientInterceptor, UnaryStreamClientInterceptor): +class RequestIdClientInterceptor( + UnaryUnaryClientInterceptor, UnaryStreamClientInterceptor +): """Injects a unique request-id header into outgoing gRPC calls if not already present.""" def __init__(self, client_version: str = "1.0.0") -> None: self.client_version = client_version - def _inject_metadata(self, client_call_details: ClientCallDetails) -> ClientCallDetails: + def _inject_metadata( + self, client_call_details: ClientCallDetails + ) -> ClientCallDetails: metadata = list(client_call_details.metadata or []) has_req_id = any(k == REQUEST_ID_HEADER for k, _ in metadata) if not has_req_id: @@ -92,7 +96,9 @@ async def intercept_service( if handler.unary_unary: unary_unary_fn = handler.unary_unary - async def logged_unary_unary(request: Any, context: grpc.aio.ServicerContext[Any, Any]) -> Any: + async def logged_unary_unary( + request: Any, context: grpc.aio.ServicerContext[Any, Any] + ) -> Any: start = time.perf_counter() peer = context.peer() try: @@ -123,7 +129,9 @@ async def logged_unary_unary(request: Any, context: grpc.aio.ServicerContext[Any duration_ms, exc, ) - await context.abort(grpc.StatusCode.INTERNAL, f"Internal server error: {exc}") + await context.abort( + grpc.StatusCode.INTERNAL, f"Internal server error: {exc}" + ) return grpc.unary_unary_rpc_method_handler( logged_unary_unary, diff --git a/src/python_grpc/core/device/pzem_004t.py b/src/python_grpc/core/device/pzem_004t.py index 09d66c8..92d39c5 100644 --- a/src/python_grpc/core/device/pzem_004t.py +++ b/src/python_grpc/core/device/pzem_004t.py @@ -10,7 +10,7 @@ import random import time -from typing import Iterator +from collections.abc import Iterator from python_grpc.proto import pzem_004t_pb2 diff --git a/tests/test_cross_host.py b/tests/test_cross_host.py index a115a1c..029cfbd 100644 --- a/tests/test_cross_host.py +++ b/tests/test_cross_host.py @@ -9,7 +9,7 @@ from __future__ import annotations import asyncio -from typing import AsyncGenerator +from collections.abc import AsyncGenerator import grpc import pytest @@ -30,8 +30,7 @@ async def cross_host_topology() -> AsyncGenerator[ DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, pzem_004t_pb2_grpc.DeviceTelemetryStub, - ], - None, + ] ]: """Spins up a central gRPC server and connects two completely independent @@ -89,7 +88,9 @@ async def test_cross_host_telemetry_streaming( received_readings: list[pzem_004t_pb2.ReadingReport] = [] async def _subscriber() -> None: - call = monitor_stub.Subscribe(pzem_004t_pb2.SubscribeRequest(device_id="PZEM-CROSS-01")) + call = monitor_stub.Subscribe( + pzem_004t_pb2.SubscribeRequest(device_id="PZEM-CROSS-01") + ) async for reading in call: received_readings.append(reading) if len(received_readings) == 3: @@ -176,7 +177,9 @@ async def test_cross_host_subscriber_disconnect_resilience( dev = PZEM004TDevice(device_id="PZEM-DEV-RESILIENT", seed=99) # 1. Connect subscriber and immediately cancel - call = monitor_stub.Subscribe(pzem_004t_pb2.SubscribeRequest(device_id="PZEM-DEV-RESILIENT")) + call = monitor_stub.Subscribe( + pzem_004t_pb2.SubscribeRequest(device_id="PZEM-DEV-RESILIENT") + ) call.cancel() # 2. Allow event loop to process cancellation diff --git a/tests/test_gateway.py b/tests/test_gateway.py index 4dcae80..6d7b533 100644 --- a/tests/test_gateway.py +++ b/tests/test_gateway.py @@ -2,7 +2,8 @@ from __future__ import annotations -from typing import AsyncGenerator +from collections.abc import AsyncGenerator + import grpc import httpx import pytest @@ -13,7 +14,9 @@ @pytest.fixture -async def gateway_client() -> AsyncGenerator[tuple[httpx.AsyncClient, DeviceTelemetryServicer], None]: +async def gateway_client() -> AsyncGenerator[ + tuple[httpx.AsyncClient, DeviceTelemetryServicer] +]: server = grpc.aio.server() servicer = DeviceTelemetryServicer() pzem_004t_pb2_grpc.add_DeviceTelemetryServicer_to_server(servicer, server) @@ -40,7 +43,7 @@ async def gateway_client() -> AsyncGenerator[tuple[httpx.AsyncClient, DeviceTele async def test_gateway_unary_post( - gateway_client: tuple[httpx.AsyncClient, DeviceTelemetryServicer] + gateway_client: tuple[httpx.AsyncClient, DeviceTelemetryServicer], ) -> None: client, servicer = gateway_client payload = { @@ -53,7 +56,7 @@ async def test_gateway_unary_post( "frequency": 50.0, "power_factor": 0.99, } - response = await client.post("/api/v1/telemetry", json=payload) + response = await client.post("/api/telemetry", json=payload) assert response.status_code == 200 data = response.json() assert data["success"] is True @@ -62,7 +65,7 @@ async def test_gateway_unary_post( async def test_gateway_batch_post( - gateway_client: tuple[httpx.AsyncClient, DeviceTelemetryServicer] + gateway_client: tuple[httpx.AsyncClient, DeviceTelemetryServicer], ) -> None: client, servicer = gateway_client readings = [ @@ -78,7 +81,7 @@ async def test_gateway_batch_post( } for i in range(3) ] - response = await client.post("/api/v1/telemetry/batch", json=readings) + response = await client.post("/api/telemetry/batch", json=readings) assert response.status_code == 200 data = response.json() assert data["received"] == 3 diff --git a/tests/test_health_and_interceptors.py b/tests/test_health_and_interceptors.py index ab5297d..ed46450 100644 --- a/tests/test_health_and_interceptors.py +++ b/tests/test_health_and_interceptors.py @@ -4,21 +4,29 @@ from __future__ import annotations -from typing import AsyncGenerator -import pytest +from collections.abc import AsyncGenerator + import grpc +import pytest from grpc_health.v1 import health, health_pb2, health_pb2_grpc from python_grpc.apps.collector.servicer import DeviceTelemetryServicer -from python_grpc.core.common.interceptors import RequestIdClientInterceptor, ServerLoggingAndRecoveryInterceptor +from python_grpc.core.common.interceptors import ( + RequestIdClientInterceptor, + ServerLoggingAndRecoveryInterceptor, +) from python_grpc.core.device.pzem_004t import PZEM004TDevice from python_grpc.proto import pzem_004t_pb2_grpc @pytest.fixture async def health_and_interceptors_env() -> AsyncGenerator[ - tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, health_pb2_grpc.HealthStub, str], - None, + tuple[ + DeviceTelemetryServicer, + pzem_004t_pb2_grpc.DeviceTelemetryStub, + health_pb2_grpc.HealthStub, + str, + ] ]: server = grpc.aio.server(interceptors=[ServerLoggingAndRecoveryInterceptor()]) servicer = DeviceTelemetryServicer() @@ -50,8 +58,11 @@ async def health_and_interceptors_env() -> AsyncGenerator[ async def test_health_check_returns_serving( health_and_interceptors_env: tuple[ - DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, health_pb2_grpc.HealthStub, str - ] + DeviceTelemetryServicer, + pzem_004t_pb2_grpc.DeviceTelemetryStub, + health_pb2_grpc.HealthStub, + str, + ], ) -> None: _, _, health_stub, service_name = health_and_interceptors_env res = await health_stub.Check(health_pb2.HealthCheckRequest(service=service_name)) @@ -60,19 +71,27 @@ async def test_health_check_returns_serving( async def test_health_check_unknown_service_returns_not_found( health_and_interceptors_env: tuple[ - DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, health_pb2_grpc.HealthStub, str - ] + DeviceTelemetryServicer, + pzem_004t_pb2_grpc.DeviceTelemetryStub, + health_pb2_grpc.HealthStub, + str, + ], ) -> None: _, _, health_stub, _ = health_and_interceptors_env with pytest.raises(grpc.aio.AioRpcError) as ctx: - await health_stub.Check(health_pb2.HealthCheckRequest(service="non.existent.Service")) + await health_stub.Check( + health_pb2.HealthCheckRequest(service="non.existent.Service") + ) assert ctx.value.code() == grpc.StatusCode.NOT_FOUND async def test_unary_call_with_client_and_server_interceptors( health_and_interceptors_env: tuple[ - DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, health_pb2_grpc.HealthStub, str - ] + DeviceTelemetryServicer, + pzem_004t_pb2_grpc.DeviceTelemetryStub, + health_pb2_grpc.HealthStub, + str, + ], ) -> None: servicer, stub, _, _ = health_and_interceptors_env device = PZEM004TDevice(seed=99) diff --git a/tests/test_mqtt_pipeline.py b/tests/test_mqtt_pipeline.py new file mode 100644 index 0000000..ef354fa --- /dev/null +++ b/tests/test_mqtt_pipeline.py @@ -0,0 +1,128 @@ +"""Integration test for MQTT-to-gRPC Ingestion Pipeline. + +Tests the simulated ESP32 MQTT Sensor Node publishing over MQTT, +the MQTT-to-gRPC Bridge receiving and transforming to Protobuf, +and the gRPC Telemetry Collector receiving and acknowledging the readings. +""" + +from __future__ import annotations + +import asyncio +from collections.abc import AsyncGenerator, Callable +from typing import Any +from unittest.mock import patch + +import grpc +import paho.mqtt.client as mqtt +import pytest + +from python_grpc.apps.collector.servicer import DeviceTelemetryServicer +from python_grpc.apps.mqtt_bridge.app import bridge_mqtt_to_grpc +from python_grpc.apps.mqtt_sensor_node.app import run_sensor_node +from python_grpc.proto import pzem_004t_pb2_grpc + + +@pytest.fixture +async def ephemeral_grpc_collector() -> AsyncGenerator[ + tuple[DeviceTelemetryServicer, str] +]: + server = grpc.aio.server() + servicer = DeviceTelemetryServicer() + pzem_004t_pb2_grpc.add_DeviceTelemetryServicer_to_server(servicer, server) + port = server.add_insecure_port("127.0.0.1:0") + await server.start() + try: + yield servicer, f"127.0.0.1:{port}" + finally: + await server.stop(grace=0) + + +async def test_mqtt_bridge_to_grpc_collector_pipeline( + ephemeral_grpc_collector: tuple[DeviceTelemetryServicer, str], +) -> None: + """Verifies MQTT messages correctly convert to ReadingReport and get stored in gRPC servicer.""" + servicer, grpc_target = ephemeral_grpc_collector + + # Registered mock clients across the test + subscribers: list[tuple[str, Any]] = [] + + class MockPahoClient: + on_connect: Callable[..., None] | None + on_message: Callable[..., None] | None + + def __init__(self, callback_api_version: mqtt.CallbackAPIVersion) -> None: + self.callback_api_version = callback_api_version + self.on_connect = None + self.on_message = None + self._is_connected = False + + def connect(self, host: str, port: int, keepalive: int = 60) -> int: + self._is_connected = True + if self.on_connect: + # Trigger successful connection callback (VERSION2 signature) + flags = mqtt.ConnectFlags(session_present=False) + rc = mqtt.ReasonCode(mqtt.PacketTypes.CONNACK, "Success") + self.on_connect(self, None, flags, rc, None) + return 0 + + def loop_start(self) -> int: + return 0 + + def loop_stop(self) -> int: + return 0 + + def disconnect(self) -> int: + self._is_connected = False + return 0 + + def subscribe(self, topic: str) -> tuple[int, int]: + subscribers.append((topic, self)) + return 0, 1 + + def publish( + self, topic: str, payload: str, qos: int = 1 + ) -> mqtt.MQTTMessageInfo: + # Deliver message to subscribers + msg = mqtt.MQTTMessage(topic=topic.encode("utf-8")) + msg.payload = payload.encode("utf-8") + msg.qos = qos + for _sub_topic, client in subscribers: + if client.on_message: + client.on_message(client, None, msg) + return mqtt.MQTTMessageInfo(0) + + with patch("paho.mqtt.client.Client", side_effect=MockPahoClient): + # 1. Run bridge in background (stop after forwarding 2 messages) + bridge_task = asyncio.create_task( + bridge_mqtt_to_grpc( + mqtt_host="mock-broker", + mqtt_port=1883, + grpc_target=grpc_target, + stop_after_count=2, + ) + ) + + # Allow bridge to start and subscribe + await asyncio.sleep(0.05) + + # 2. Run simulated sensor node in a background thread because run_sensor_node is synchronous + loop = asyncio.get_running_loop() + sensor_task = loop.run_in_executor( + None, + run_sensor_node, + "mock-broker", + 1883, + "ESP32-UNIT-TEST", + "devices", + 0.01, + 2, + ) + + await asyncio.gather(bridge_task, sensor_task) + + # 3. Assert gRPC collector received both readings + assert len(servicer.readings) == 2 + assert servicer.readings[0].device_id == "ESP32-UNIT-TEST" + assert servicer.readings[1].device_id == "ESP32-UNIT-TEST" + assert servicer.readings[0].voltage > 0.0 + assert servicer.readings[0].current > 0.0 diff --git a/tests/test_telemetry.py b/tests/test_telemetry.py index ad93b72..232e919 100644 --- a/tests/test_telemetry.py +++ b/tests/test_telemetry.py @@ -3,7 +3,8 @@ from __future__ import annotations import asyncio -from typing import AsyncGenerator +from collections.abc import AsyncGenerator + import grpc import pytest @@ -14,8 +15,9 @@ @pytest.fixture async def telemetry_env() -> AsyncGenerator[ - tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice], - None, + tuple[ + DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice + ] ]: server = grpc.aio.server() servicer = DeviceTelemetryServicer() @@ -35,7 +37,9 @@ async def telemetry_env() -> AsyncGenerator[ async def test_valid_reading_is_acknowledged_and_stored( - telemetry_env: tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice] + telemetry_env: tuple[ + DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice + ], ) -> None: servicer, stub, device = telemetry_env ack = await stub.ReportReading(device.read()) @@ -46,7 +50,9 @@ async def test_valid_reading_is_acknowledged_and_stored( async def test_invalid_reading_is_rejected( - telemetry_env: tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice] + telemetry_env: tuple[ + DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice + ], ) -> None: servicer, stub, device = telemetry_env report = device.read() @@ -57,7 +63,9 @@ async def test_invalid_reading_is_rejected( async def test_batch_upload_produces_summary( - telemetry_env: tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice] + telemetry_env: tuple[ + DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice + ], ) -> None: servicer, stub, device = telemetry_env @@ -79,7 +87,9 @@ async def readings(): async def test_streams_live_readings_for_device( - telemetry_env: tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice] + telemetry_env: tuple[ + DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice + ], ) -> None: _, stub, device = telemetry_env request = pzem_004t_pb2.SubscribeRequest(device_id=device.device_id) @@ -113,7 +123,9 @@ async def _producer() -> None: async def test_every_reading_gets_an_ack( - telemetry_env: tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice] + telemetry_env: tuple[ + DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, PZEM004TDevice + ], ) -> None: servicer, stub, device = telemetry_env