Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 9 additions & 9 deletions README.md
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -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.
Expand Down
57 changes: 38 additions & 19 deletions scripts/gen_proto.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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__":
Expand Down
126 changes: 0 additions & 126 deletions scripts/simulate_cross_host.py

This file was deleted.

4 changes: 2 additions & 2 deletions src/python_grpc/apps/collector/app.py
Original file line number Diff line number Diff line change
@@ -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.
"""
Expand Down Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion src/python_grpc/core/common/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading