From 47a7962fa3e4fae70154876d03b2a11eaa6de66b Mon Sep 17 00:00:00 2001 From: Farrel Augusta Dinata Date: Mon, 14 Sep 2026 08:25:39 +0700 Subject: [PATCH 1/2] refactor: remove `mqtt-bridge` service --- README.md | 38 ++-- docker-compose.yml | 25 +-- src/python_grpc/apps/collector/app.py | 50 ++++- .../apps/collector/mqtt_consumer.py | 106 ++++++++++ src/python_grpc/apps/collector/servicer.py | 20 ++ src/python_grpc/apps/mqtt_bridge/__init__.py | 5 - src/python_grpc/apps/mqtt_bridge/__main__.py | 6 - src/python_grpc/apps/mqtt_bridge/app.py | 181 ------------------ .../apps/rest_gateway/routers/telemetry.py | 35 ++++ tests/test_mqtt_pipeline.py | 81 ++++---- 10 files changed, 279 insertions(+), 268 deletions(-) create mode 100644 src/python_grpc/apps/collector/mqtt_consumer.py delete mode 100644 src/python_grpc/apps/mqtt_bridge/__init__.py delete mode 100644 src/python_grpc/apps/mqtt_bridge/__main__.py delete mode 100644 src/python_grpc/apps/mqtt_bridge/app.py diff --git a/README.md b/README.md index af796c5..1e6a42e 100644 --- a/README.md +++ b/README.md @@ -19,12 +19,12 @@ flowchart LR Broker["Eclipse Mosquitto\n(devices/+/telemetry)"] end - subgraph Bridges["Ingestion Bridges"] - Bridge["MQTT-to-gRPC Bridge\n(apps/mqtt_bridge)\n[paho-mqtt + grpc.aio]"] - end - subgraph Collector["Telemetry Collector Service (Port 50051)"] - Coll["Collector Service\n(apps/collector)\n(gRPC Ingestion & Pub/Sub Hub)"] + MQTTSub["Embedded MQTT Consumer\n(paho-mqtt)"] + Serv["DeviceTelemetryServicer\n(Pub/Sub & In-Memory Store)"] + gRPCServer["gRPC Server (HTTP/2)\n(Port 50051)"] + MQTTSub -->|"Internal Ingest"| Serv + Serv <--> gRPCServer end subgraph RESTConsumer["Web / Mobile / Dashboard"] @@ -32,23 +32,22 @@ flowchart LR end ESP32 -->|"MQTT Publish (JSON)"| Broker - Broker -->|"MQTT Subscribe"| Bridge - Bridge ==>|"gRPC ReportReading (HTTP/2)"| Coll - REST ==>|"gRPC Unary & Batch (HTTP/2)"| Coll + Broker -->|"MQTT Subscribe"| MQTTSub + REST ==>|"gRPC Unary & Batch (HTTP/2)"| gRPCServer ``` ### Industry-Grade Capabilities Implemented - **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. + - **Embedded MQTT Ingestion**: Telemetry Collector directly consumes MQTT telemetry topics into its real-time pub/sub hub. - **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. -- **Container Orchestration**: Multi-container Docker Compose topology orchestrating Mosquitto MQTT, Collector, MQTT Bridge, and REST Gateway across segmented bridge networks. +- **Container Orchestration**: Multi-container Docker Compose topology orchestrating Mosquitto MQTT, Collector (with embedded MQTT subscriber), and REST Gateway across segmented bridge networks. ### The four gRPC call types @@ -74,16 +73,14 @@ src/python_grpc/ device/ pzem_004t.py # PZEM004TDevice hardware physics simulator apps/ # autonomous deployable applications - collector/ # Cloud-tier: Telemetry Collector Server (gRPC only) - app.py # server lifecycle, health check, graceful shutdown + collector/ # Cloud-tier: Telemetry Collector Server (gRPC + Embedded MQTT) + app.py # server lifecycle, embedded MQTT consumer, health check + mqtt_consumer.py # paho-mqtt background subscriber feeding servicer directly servicer.py # in-memory pub/sub telemetry broadcast servicer __main__.py # CLI entry point (python -m python_grpc.apps.collector) 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 @@ -100,7 +97,7 @@ 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 + test_mqtt_pipeline.py # Embedded MQTT ingestion & gRPC live streaming tests Dockerfile # multi-app container build docker-compose.yml # multi-network orchestration with Mosquitto MQTT broker ``` @@ -123,7 +120,7 @@ poetry install #### 1. Start the Telemetry Collector Server (Terminal 1) ```bash -poetry run python -m python_grpc.apps.collector --host 0.0.0.0 --port 50051 +poetry run python -m python_grpc.apps.collector --host 0.0.0.0 --port 50051 --mqtt-host localhost --mqtt-port 1883 ``` #### 2. Start the REST-to-gRPC Gateway (Terminal 2) @@ -137,18 +134,15 @@ curl -X POST "http://localhost:8000/api/telemetry" \ -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) +#### 3. Run the Simulated ESP32 Sensor Node (Terminal 3) 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 (Mosquitto MQTT broker, Collector, MQTT Bridge, and REST Gateway): +Spin up the entire microservice topology (Mosquitto MQTT broker, Collector with embedded MQTT, and REST Gateway): ```bash docker compose up --build ``` diff --git a/docker-compose.yml b/docker-compose.yml index c9e1224..0455866 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -7,15 +7,20 @@ services: build: context: . dockerfile: Dockerfile - command: python -m python_grpc.apps.collector --host 0.0.0.0 --port 50051 + command: python -m python_grpc.apps.collector --host 0.0.0.0 --port 50051 --mqtt-host broker.internal --mqtt-port 1883 ports: - "50051:50051" + depends_on: + - mqtt-broker environment: - LOG_LEVEL=${LOG_LEVEL:-INFO} + - MQTT_BROKER_HOST=broker.internal + - MQTT_BROKER_PORT=1883 networks: cloud-tier: aliases: - collector.internal + edge-tier: restart: unless-stopped @@ -70,24 +75,6 @@ services: 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: - - LOG_LEVEL=${LOG_LEVEL:-INFO} - networks: - edge-tier: - cloud-tier: - restart: unless-stopped networks: cloud-tier: diff --git a/src/python_grpc/apps/collector/app.py b/src/python_grpc/apps/collector/app.py index 076be8d..31a21fa 100644 --- a/src/python_grpc/apps/collector/app.py +++ b/src/python_grpc/apps/collector/app.py @@ -15,6 +15,7 @@ import grpc from grpc_health.v1 import health, health_pb2, health_pb2_grpc +from python_grpc.apps.collector.mqtt_consumer import CollectorMQTTConsumer from python_grpc.apps.collector.servicer import DeviceTelemetryServicer from python_grpc.core.common.config import DEFAULT_GRPC_CHANNEL_OPTIONS from python_grpc.core.common.interceptors import ServerLoggingAndRecoveryInterceptor @@ -23,7 +24,13 @@ logger = logging.getLogger("telemetry_collector_server") -async def serve(host: str = "0.0.0.0", port: int = 50051) -> None: +async def serve( + host: str = "0.0.0.0", + port: int = 50051, + mqtt_host: str | None = None, + mqtt_port: int = 1883, + mqtt_topic: str = "devices/+/telemetry", +) -> None: # 1. Initialize server with production channel options and interceptors interceptors = [ServerLoggingAndRecoveryInterceptor()] server = grpc.aio.server( @@ -49,7 +56,18 @@ async def serve(host: str = "0.0.0.0", port: int = 50051) -> None: await server.start() logger.info("Collector service started on %s (Health check enabled)", listen_addr) - # 4. Graceful shutdown handling + # 4. Start embedded MQTT consumer if host is configured + mqtt_consumer: CollectorMQTTConsumer | None = None + if mqtt_host: + mqtt_consumer = CollectorMQTTConsumer( + servicer=servicer, + broker_host=mqtt_host, + broker_port=mqtt_port, + topic=mqtt_topic, + ) + mqtt_consumer.start() + + # 5. Graceful shutdown handling stop_event = asyncio.Event() loop = asyncio.get_running_loop() @@ -79,6 +97,8 @@ def _trigger_stop(*args: Any) -> None: for task in pending: task.cancel() finally: + if mqtt_consumer is not None: + mqtt_consumer.stop() logger.info("Stopping gRPC server with 5s grace period...") await server.stop(grace=5.0) logger.info("Collector service stopped cleanly.") @@ -99,6 +119,22 @@ def main() -> None: default=int(os.getenv("COLLECTOR_PORT", "50051")), help="Port to listen on", ) + parser.add_argument( + "--mqtt-host", + default=os.getenv("MQTT_BROKER_HOST", None), + help="Optional MQTT broker host to subscribe to sensor readings", + ) + parser.add_argument( + "--mqtt-port", + type=int, + default=int(os.getenv("MQTT_BROKER_PORT", "1883")), + help="MQTT broker port", + ) + parser.add_argument( + "--mqtt-topic", + default=os.getenv("MQTT_TOPIC", "devices/+/telemetry"), + help="MQTT topic filter to subscribe to", + ) parser.add_argument( "--log-level", default=os.getenv("LOG_LEVEL", "INFO"), @@ -112,7 +148,15 @@ def main() -> None: format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) try: - asyncio.run(serve(host=args.host, port=args.port)) + asyncio.run( + serve( + host=args.host, + port=args.port, + mqtt_host=args.mqtt_host, + mqtt_port=args.mqtt_port, + mqtt_topic=args.mqtt_topic, + ) + ) except KeyboardInterrupt: logger.info("Application interrupted by user.") diff --git a/src/python_grpc/apps/collector/mqtt_consumer.py b/src/python_grpc/apps/collector/mqtt_consumer.py new file mode 100644 index 0000000..b99000f --- /dev/null +++ b/src/python_grpc/apps/collector/mqtt_consumer.py @@ -0,0 +1,106 @@ +"""Embedded MQTT consumer for the Telemetry Collector application. + +Subscribes to telemetry topics on an MQTT broker and directly feeds received +PZEM-004T readings into the collector servicer and its pub/sub broadcast hub. +""" + +from __future__ import annotations + +import json +import logging + +import paho.mqtt.client as mqtt + +from python_grpc.apps.collector.servicer import DeviceTelemetryServicer +from python_grpc.proto import pzem_004t_pb2 + +logger = logging.getLogger("collector_mqtt_consumer") + + +class CollectorMQTTConsumer: + """Manages an embedded MQTT subscriber client feeding readings directly into servicer.""" + + def __init__( + self, + servicer: DeviceTelemetryServicer, + broker_host: str, + broker_port: int = 1883, + topic: str = "devices/+/telemetry", + ) -> None: + self.servicer = servicer + self.broker_host = broker_host + self.broker_port = broker_port + self.topic = topic + self._client: mqtt.Client | None = None + + def start(self) -> None: + """Connect to the MQTT broker and start background listening.""" + logger.info( + "Starting embedded MQTT consumer: %s:%d (topic=%r)", + self.broker_host, + self.broker_port, + self.topic, + ) + self._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...", self.topic) + client.subscribe(self.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)), + ) + ack = self.servicer.ingest_reading(reading) + logger.debug("MQTT reading ingested: %s", ack.message) + except (json.JSONDecodeError, KeyError, ValueError) as err: + logger.warning("Failed to parse MQTT message payload: %s", err) + + self._client.on_connect = on_connect + self._client.on_message = on_message + + try: + self._client.connect(self.broker_host, self.broker_port, keepalive=60) + self._client.loop_start() + logger.info("Embedded MQTT consumer loop started.") + except Exception as exc: + logger.warning( + "Could not connect to MQTT broker at %s:%d (%s). Operating in gRPC-only mode.", + self.broker_host, + self.broker_port, + exc, + ) + + def stop(self) -> None: + """Cleanly stop MQTT client loop and disconnect.""" + if self._client is not None: + logger.info("Stopping embedded MQTT consumer...") + self._client.loop_stop() + self._client.disconnect() + self._client = None + logger.info("Embedded MQTT consumer stopped.") diff --git a/src/python_grpc/apps/collector/servicer.py b/src/python_grpc/apps/collector/servicer.py index 953fd63..8b05040 100644 --- a/src/python_grpc/apps/collector/servicer.py +++ b/src/python_grpc/apps/collector/servicer.py @@ -79,6 +79,26 @@ def _broadcast(self, report: pzem_004t_pb2.ReadingReport) -> None: "Subscriber queue full, dropping reading for %s", report.device_id ) + def ingest_reading(self, report: pzem_004t_pb2.ReadingReport) -> pzem_004t_pb2.Ack: + """Ingest a reading from an internal/embedded source (such as MQTT) and broadcast.""" + if not _is_valid(report): + logger.warning( + "Rejected invalid internal reading from %s", report.device_id + ) + return _ack(False, f"rejected invalid reading from {report.device_id}") + self._readings.append(report) + self._broadcast(report) + logger.info( + "Ingested reading device=%s v=%.1fV i=%.2fA p=%.1fW", + report.device_id, + report.voltage, + report.current, + report.active_power, + ) + return _ack( + True, f"stored reading #{len(self._readings)} from {report.device_id}" + ) + @override async def ReportReading( self, diff --git a/src/python_grpc/apps/mqtt_bridge/__init__.py b/src/python_grpc/apps/mqtt_bridge/__init__.py deleted file mode 100644 index c3572c3..0000000 --- a/src/python_grpc/apps/mqtt_bridge/__init__.py +++ /dev/null @@ -1,5 +0,0 @@ -"""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 deleted file mode 100644 index 354e42c..0000000 --- a/src/python_grpc/apps/mqtt_bridge/__main__.py +++ /dev/null @@ -1,6 +0,0 @@ -"""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 deleted file mode 100644 index 5bb609f..0000000 --- a/src/python_grpc/apps/mqtt_bridge/app.py +++ /dev/null @@ -1,181 +0,0 @@ -"""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/rest_gateway/routers/telemetry.py b/src/python_grpc/apps/rest_gateway/routers/telemetry.py index 68d63ae..08caac7 100644 --- a/src/python_grpc/apps/rest_gateway/routers/telemetry.py +++ b/src/python_grpc/apps/rest_gateway/routers/telemetry.py @@ -2,11 +2,14 @@ from __future__ import annotations +import asyncio +import json from collections.abc import AsyncGenerator from typing import Annotated import grpc from fastapi import APIRouter, Depends, HTTPException +from fastapi.responses import StreamingResponse from python_grpc.apps.rest_gateway.dependencies import get_telemetry_stub from python_grpc.apps.rest_gateway.schemas import ( @@ -85,3 +88,35 @@ async def _generator() -> AsyncGenerator[pzem_004t_pb2.ReadingReport]: raise HTTPException( status_code=502, detail=f"gRPC call failed: {exc.details()}" ) from exc + + +@router.get("/live") +async def stream_live_telemetry( + stub: TelemetryStubDep, + device_id: str = "*", +) -> StreamingResponse: + """Stream live telemetry from gRPC Collector to HTTP clients via Server-Sent Events (SSE).""" + + async def _event_generator() -> AsyncGenerator[str]: + request = pzem_004t_pb2.SubscribeRequest(device_id=device_id) + call = stub.Subscribe(request) + try: + async for reading in call: + 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, + } + yield f"data: {json.dumps(payload)}\n\n" + except (grpc.RpcError, asyncio.CancelledError): + return + finally: + call.cancel() + + return StreamingResponse(_event_generator(), media_type="text/event-stream") diff --git a/tests/test_mqtt_pipeline.py b/tests/test_mqtt_pipeline.py index ef354fa..aa65017 100644 --- a/tests/test_mqtt_pipeline.py +++ b/tests/test_mqtt_pipeline.py @@ -1,8 +1,9 @@ -"""Integration test for MQTT-to-gRPC Ingestion Pipeline. +"""Integration test for Embedded MQTT Ingestion in Collector & gRPC Streaming. 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. +the Telemetry Collector ingesting readings via its embedded MQTT consumer, +and a gRPC client (such as the REST Gateway or monitoring station) +receiving live telemetry streamed via gRPC Subscribe(). """ from __future__ import annotations @@ -16,34 +17,39 @@ import paho.mqtt.client as mqtt import pytest +from python_grpc.apps.collector.mqtt_consumer import CollectorMQTTConsumer 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 +from python_grpc.proto import pzem_004t_pb2, pzem_004t_pb2_grpc @pytest.fixture -async def ephemeral_grpc_collector() -> AsyncGenerator[ - tuple[DeviceTelemetryServicer, str] +async def collector_with_mqtt() -> AsyncGenerator[ + tuple[DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, 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() + + channel = grpc.aio.insecure_channel(f"127.0.0.1:{port}") + stub = pzem_004t_pb2_grpc.DeviceTelemetryStub(channel) try: - yield servicer, f"127.0.0.1:{port}" + yield servicer, stub, f"127.0.0.1:{port}" finally: + await channel.close() await server.stop(grace=0) -async def test_mqtt_bridge_to_grpc_collector_pipeline( - ephemeral_grpc_collector: tuple[DeviceTelemetryServicer, str], +async def test_embedded_mqtt_collector_to_grpc_streaming( + collector_with_mqtt: tuple[ + DeviceTelemetryServicer, pzem_004t_pb2_grpc.DeviceTelemetryStub, str + ], ) -> None: - """Verifies MQTT messages correctly convert to ReadingReport and get stored in gRPC servicer.""" - servicer, grpc_target = ephemeral_grpc_collector + """Verifies that MQTT messages are ingested directly into the Collector and streamed out via gRPC Subscribe.""" + servicer, stub, _target = collector_with_mqtt - # Registered mock clients across the test subscribers: list[tuple[str, Any]] = [] class MockPahoClient: @@ -59,7 +65,6 @@ def __init__(self, callback_api_version: mqtt.CallbackAPIVersion) -> None: 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) @@ -82,7 +87,6 @@ def subscribe(self, topic: str) -> tuple[int, int]: 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 @@ -92,37 +96,50 @@ def publish( 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, - ) + # 1. Start embedded MQTT consumer in the collector + mqtt_consumer = CollectorMQTTConsumer( + servicer=servicer, + broker_host="mock-broker", + broker_port=1883, ) + mqtt_consumer.start() + + # 2. Open gRPC live subscriber stream (simulating REST Gateway / dashboard) + received_over_grpc: list[pzem_004t_pb2.ReadingReport] = [] - # Allow bridge to start and subscribe + async def _grpc_subscriber() -> None: + call = stub.Subscribe(pzem_004t_pb2.SubscribeRequest(device_id="*")) + async for reading in call: + received_over_grpc.append(reading) + if len(received_over_grpc) == 2: + call.cancel() + break + + subscriber_task = asyncio.create_task(_grpc_subscriber()) await asyncio.sleep(0.05) - # 2. Run simulated sensor node in a background thread because run_sensor_node is synchronous + # 3. Run simulated ESP32 sensor node publishing over MQTT loop = asyncio.get_running_loop() sensor_task = loop.run_in_executor( None, run_sensor_node, "mock-broker", 1883, - "ESP32-UNIT-TEST", + "ESP32-EMBEDDED-01", "devices", 0.01, 2, ) - await asyncio.gather(bridge_task, sensor_task) + await asyncio.gather(subscriber_task, sensor_task) + mqtt_consumer.stop() - # 3. Assert gRPC collector received both readings + # 4. Assert both readings were ingested by the servicer 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 + assert servicer.readings[0].device_id == "ESP32-EMBEDDED-01" + + # 5. Assert readings were streamed over gRPC in real time to the subscriber + assert len(received_over_grpc) == 2 + assert received_over_grpc[0].device_id == "ESP32-EMBEDDED-01" + assert received_over_grpc[0].voltage > 0.0 + assert received_over_grpc[1].voltage > 0.0 From 484c64385e416d5dbda11f290919910dd9c6e280 Mon Sep 17 00:00:00 2001 From: Farrel Augusta Dinata Date: Mon, 14 Sep 2026 08:35:01 +0700 Subject: [PATCH 2/2] refactor: enhance Dockerfile and entrypoint script for improved service management --- Dockerfile | 49 +++++++++++++++++++++++++++++++++++------- docker-compose.yml | 6 +++--- entrypoint.sh | 53 ++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 97 insertions(+), 11 deletions(-) create mode 100644 entrypoint.sh diff --git a/Dockerfile b/Dockerfile index a7dc44e..8acd7e0 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,8 +1,11 @@ -FROM python:3.13-slim +# ============================================================= +# Stage 1: Build & Dependency Resolution Stage +# ============================================================= +FROM python:3.13-slim AS builder -WORKDIR /app +WORKDIR /build -# Install system dependencies +# Install build dependencies RUN apt-get update && apt-get install -y --no-install-recommends \ curl \ build-essential \ @@ -11,13 +14,43 @@ RUN apt-get update && apt-get install -y --no-install-recommends \ # Install poetry RUN pip install --no-cache-dir poetry -COPY pyproject.toml poetry.lock* README.md /app/ -RUN poetry config virtualenvs.create false \ - && poetry install --only main --no-interaction --no-ansi +# Copy dependency manifests +COPY pyproject.toml poetry.lock* README.md /build/ + +# Create a standalone virtual environment in /build/.venv +RUN poetry config virtualenvs.in-project true \ + && poetry install --only main --no-interaction --no-ansi --no-root + +# ============================================================= +# Stage 2: Final Lean Production Runtime Stage +# ============================================================= +FROM python:3.13-slim AS runtime + +WORKDIR /app + +# Create a non-root application user for container security +RUN groupadd -r appuser && useradd -r -g appuser -d /app -s /sbin/nologin appuser + +# Copy virtual environment from builder stage +COPY --from=builder /build/.venv /app/.venv +# Copy application source code and entrypoint script COPY src /app/src -ENV PYTHONPATH=/app/src +COPY entrypoint.sh /app/entrypoint.sh +RUN chmod +x /app/entrypoint.sh + +# Set environment variables +ENV PATH="/app/.venv/bin:$PATH" \ + PYTHONPATH="/app/src" \ + PYTHONUNBUFFERED=1 \ + PYTHONDONTWRITEBYTECODE=1 + +# Change file ownership to non-root user +RUN chown -R appuser:appuser /app + +USER appuser EXPOSE 50051 8000 -CMD ["python", "-m", "python_grpc.apps.collector"] +ENTRYPOINT ["/app/entrypoint.sh"] +CMD ["help"] diff --git a/docker-compose.yml b/docker-compose.yml index 0455866..c1f014a 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -7,7 +7,7 @@ services: build: context: . dockerfile: Dockerfile - command: python -m python_grpc.apps.collector --host 0.0.0.0 --port 50051 --mqtt-host broker.internal --mqtt-port 1883 + command: collector --host 0.0.0.0 --port 50051 --mqtt-host broker.internal --mqtt-port 1883 ports: - "50051:50051" depends_on: @@ -32,7 +32,7 @@ services: build: context: . dockerfile: Dockerfile - command: python -m python_grpc.apps.rest_gateway --host 0.0.0.0 --port 8000 --grpc-target collector.internal:50051 + command: gateway --host 0.0.0.0 --port 8000 --grpc-target collector.internal:50051 ports: - "8000:8000" depends_on: @@ -66,7 +66,7 @@ services: build: context: . dockerfile: Dockerfile - command: python -m python_grpc.apps.mqtt_sensor_node --broker-host broker.internal --broker-port 1883 --device-id ESP32-PZEM-01 --count 5 + command: sensor-node --broker-host broker.internal --broker-port 1883 --device-id ESP32-PZEM-01 --count 5 depends_on: - mqtt-broker environment: diff --git a/entrypoint.sh b/entrypoint.sh new file mode 100644 index 0000000..d5c266f --- /dev/null +++ b/entrypoint.sh @@ -0,0 +1,53 @@ +#!/bin/sh +set -e + +# ============================================================================== +# Unified Container Entrypoint Script +# Dispatches execution based on service name or executes custom commands. +# ============================================================================== + +show_help() { + echo "Python gRPC & MQTT IoT Multi-Service Container" + echo "" + echo "Usage: docker run [service_name|command] [args...]" + echo "" + echo "Available Services:" + echo " collector Run central gRPC collector & embedded MQTT consumer" + echo " gateway Run FastAPI REST-to-gRPC gateway" + echo " sensor-node Run simulated MQTT sensor node (ESP32 PZEM-004T)" + echo " help Show this help message" + echo "" + echo "Examples:" + echo " docker run collector --host 0.0.0.0 --port 50051" + echo " docker run gateway --host 0.0.0.0 --port 8000" + echo " docker run sensor-node --broker-host localhost --count 10" + echo " docker run python -m pytest tests/" +} + +case "$1" in + collector) + shift + exec python -m python_grpc.apps.collector "$@" + ;; + gateway|rest-gateway) + shift + exec python -m python_grpc.apps.rest_gateway "$@" + ;; + sensor-node|mqtt-sensor-node) + shift + exec python -m python_grpc.apps.mqtt_sensor_node "$@" + ;; + help|--help|-h) + show_help + exit 0 + ;; + "") + echo "No command specified. Displaying usage:" + show_help + exit 1 + ;; + *) + # Allow passing arbitrary commands (e.g., bash, pytest, python, etc.) + exec "$@" + ;; +esac