From 92d00c1246082857af93e88f6ca7dd994eb72a9a Mon Sep 17 00:00:00 2001 From: Farrel Augusta Dinata Date: Mon, 14 Sep 2026 08:43:14 +0700 Subject: [PATCH 1/2] docs: update README for clarity and adjust feature descriptions --- README.md | 18 +++++++++--------- src/python_grpc/apps/collector/app.py | 4 ++-- src/python_grpc/core/common/config.py | 2 +- 3 files changed, 12 insertions(+), 12 deletions(-) diff --git a/README.md b/README.md index 1e6a42e..9a60a43 100644 --- a/README.md +++ b/README.md @@ -1,11 +1,11 @@ -# Python gRPC — PZEM-004t Industry-Grade Microservices & Gateway Prototype +# Python gRPC — PZEM-004t Experimental Microservices & Learning Lab -An end-to-end demonstration of **industry-grade gRPC in Python** using the async -`grpc.aio` API. A simulated electrical device reads live data from a -**PZEM-004t** AC power meter (voltage, current, active power, energy, -frequency, power factor) and pushes it over gRPC to a collector server, with an -optional **FastAPI REST-to-gRPC Gateway** allowing external REST API consumers to -interact with the gRPC microservice. +A hands-on learning and experimental project exploring **gRPC in Python** using the async +`grpc.aio` API. Built to understand how gRPC communication patterns work in practice, +this project simulates an electrical device reading data from a **PZEM-004t** AC power +meter (voltage, current, active power, energy, frequency, power factor) and streaming +it over gRPC to a collector server, complemented by a **FastAPI REST-to-gRPC Gateway** +to explore how external REST clients interact with internal gRPC microservices. ## Architecture @@ -36,9 +36,9 @@ flowchart LR REST ==>|"gRPC Unary & Batch (HTTP/2)"| gRPCServer ``` -### Industry-Grade Capabilities Implemented +### Key Features & Patterns Explored -- **Industrial IoT Protocol Hierarchy**: +- **IoT Protocol Integration**: - **MQTT**: Lightweight pub/sub for resource-constrained microcontrollers (ESP32) reading PZEM-004T sensors. - **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. diff --git a/src/python_grpc/apps/collector/app.py b/src/python_grpc/apps/collector/app.py index 31a21fa..566b1e8 100644 --- a/src/python_grpc/apps/collector/app.py +++ b/src/python_grpc/apps/collector/app.py @@ -1,4 +1,4 @@ -"""Async gRPC telemetry collector application bootstrap with production interceptors, +"""Async gRPC telemetry collector application bootstrap with custom interceptors, health checking, and graceful shutdown. """ @@ -31,7 +31,7 @@ async def serve( mqtt_port: int = 1883, mqtt_topic: str = "devices/+/telemetry", ) -> None: - # 1. Initialize server with production channel options and interceptors + # 1. Initialize server with channel options and interceptors interceptors = [ServerLoggingAndRecoveryInterceptor()] server = grpc.aio.server( interceptors=interceptors, diff --git a/src/python_grpc/core/common/config.py b/src/python_grpc/core/common/config.py index 70f6197..025a069 100644 --- a/src/python_grpc/core/common/config.py +++ b/src/python_grpc/core/common/config.py @@ -4,7 +4,7 @@ from typing import Any -# Production HTTP/2 channel options for keepalive and resiliency +# Recommended HTTP/2 channel options for keepalive and connection 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 From cd8db172db7da46384c0397a63df04bfc7375e40 Mon Sep 17 00:00:00 2001 From: Farrel Augusta Dinata Date: Mon, 14 Sep 2026 08:52:40 +0700 Subject: [PATCH 2/2] refactor: update gen_proto script to compile all .proto files and remove simulate_cross_host script --- scripts/gen_proto.py | 57 ++++++++++----- scripts/simulate_cross_host.py | 126 --------------------------------- 2 files changed, 38 insertions(+), 145 deletions(-) delete mode 100644 scripts/simulate_cross_host.py diff --git a/scripts/gen_proto.py b/scripts/gen_proto.py index 84519a1..29f4555 100644 --- a/scripts/gen_proto.py +++ b/scripts/gen_proto.py @@ -1,7 +1,13 @@ -"""Generate gRPC stubs from the pzem_004t.proto definition. +"""Generate gRPC stubs from all .proto definitions in the repository. -Run from the repository root (or via `poetry run python scripts/gen_proto.py`). -Writes pzem_004t_pb2.py / pzem_004t_pb2_grpc.py / .pyi next to the proto. +Run from the repository root: + poetry run python scripts/gen_proto.py + +Automatically discovers all `*.proto` files located in `src/python_grpc/proto/` +and compiles them into: + - `*_pb2.py` (Protobuf message classes) + - `*_pb2.pyi` (Protobuf type stubs) + - `*_pb2_grpc.py` (gRPC client & servicer classes) """ from __future__ import annotations @@ -12,28 +18,41 @@ ROOT = Path(__file__).resolve().parents[1] SRC = ROOT / "src" -PROTO = SRC / "python_grpc" / "proto" / "pzem_004t.proto" +PROTO_DIR = SRC / "python_grpc" / "proto" def main() -> None: - result = subprocess.run( - [ - sys.executable, - "-m", - "grpc_tools.protoc", - f"--proto_path={SRC}", - f"--python_out={SRC}", - f"--pyi_out={SRC}", - f"--grpc_python_out={SRC}", - str(PROTO), - ], - capture_output=True, - text=True, - ) + if not PROTO_DIR.exists(): + print(f"Error: Proto directory '{PROTO_DIR}' does not exist.", file=sys.stderr) + sys.exit(1) + + proto_files = sorted(PROTO_DIR.glob("*.proto")) + if not proto_files: + print(f"No .proto files found in {PROTO_DIR}", file=sys.stderr) + return + + print(f"Discovered {len(proto_files)} .proto file(s) in {PROTO_DIR}:") + for proto in proto_files: + print(f" - {proto.name}") + + cmd = [ + sys.executable, + "-m", + "grpc_tools.protoc", + f"--proto_path={SRC}", + f"--python_out={SRC}", + f"--pyi_out={SRC}", + f"--grpc_python_out={SRC}", + *[str(p) for p in proto_files], + ] + + result = subprocess.run(cmd, capture_output=True, text=True) if result.returncode != 0: + print("Protoc compilation failed:", file=sys.stderr) print(result.stderr, file=sys.stderr) raise SystemExit(result.returncode) - print("Generated stubs for", PROTO.name, "->", SRC / "python_grpc" / "proto") + + print("\nSuccessfully compiled all proto definitions into Python stubs.") if __name__ == "__main__": diff --git a/scripts/simulate_cross_host.py b/scripts/simulate_cross_host.py deleted file mode 100644 index 0ed9316..0000000 --- a/scripts/simulate_cross_host.py +++ /dev/null @@ -1,126 +0,0 @@ -"""Cross-host gRPC communication simulator. - -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 (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. -""" - -from __future__ import annotations - -import asyncio -import logging - -import grpc - -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 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" -) -logger = logging.getLogger("cross_host_sim") - - -async def main() -> None: - print("\n" + "=" * 70) - print(" SIMULATING MULTI-HOST / MULTI-SERVER gRPC TELEMETRY COMMUNICATON") - print("=" * 70 + "\n") - - # 1. Spawn Host 1 (Cloud Server) - server_logger = logging.getLogger("Host1-CloudServer") - server = grpc.aio.server(options=DEFAULT_GRPC_CHANNEL_OPTIONS) - 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() - server_logger.info("Host 1 (Collector) listening on 127.0.0.1:%d", port) - - target = f"127.0.0.1:{port}" - - # 2. Spawn Host 2 (Remote Dashboard / Monitoring Station) - host2_logger = logging.getLogger("Host2-RemoteMonitor") - received_readings: list[pzem_004t_pb2.ReadingReport] = [] - - async def host2_subscriber() -> None: - async with grpc.aio.insecure_channel( - target, - options=DEFAULT_GRPC_CHANNEL_OPTIONS, - 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") - ) - try: - async for reading in call: - host2_logger.info( - "-> [LIVE EVENT RECEIVED] Host 2 received reading from %s: %.1fV, %.2fA, %.1fW", - reading.device_id, - reading.voltage, - reading.current, - reading.active_power, - ) - received_readings.append(reading) - if len(received_readings) >= 3: - call.cancel() - break - except asyncio.CancelledError: - pass - - # 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 - device = PZEM004TDevice(device_id="PZEM-REMOTE-01", seed=123) - - async with grpc.aio.insecure_channel( - target, - options=DEFAULT_GRPC_CHANNEL_OPTIONS, - 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) - - for i in range(1, 4): - reading = device.read() - host3_logger.info( - "<- [TRANSMITTING] Host 3 sending Reading #%d (%.1fV, %.2fA)...", - i, - reading.voltage, - reading.current, - ) - ack = await stub.ReportReading(reading, timeout=5.0) - 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 - t2 = asyncio.create_task(host2_subscriber()) - t3 = asyncio.create_task(host3_edge_device()) - - 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("-" * 70 + "\n") - - await server.stop(grace=0) - - -if __name__ == "__main__": - asyncio.run(main())