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
49 changes: 41 additions & 8 deletions Dockerfile
Original file line number Diff line number Diff line change
@@ -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 \
Expand All @@ -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"]
38 changes: 16 additions & 22 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,36 +19,35 @@ 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"]
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
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

Expand All @@ -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
Expand All @@ -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
```
Expand All @@ -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)
Expand All @@ -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
```
Expand Down
29 changes: 8 additions & 21 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,15 +7,20 @@ services:
build:
context: .
dockerfile: Dockerfile
command: python -m python_grpc.apps.collector --host 0.0.0.0 --port 50051
command: 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


Expand All @@ -27,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:
Expand Down Expand Up @@ -61,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:
Expand All @@ -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:
Expand Down
53 changes: 53 additions & 0 deletions entrypoint.sh
Original file line number Diff line number Diff line change
@@ -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 <image> [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 <image> collector --host 0.0.0.0 --port 50051"
echo " docker run <image> gateway --host 0.0.0.0 --port 8000"
echo " docker run <image> sensor-node --broker-host localhost --count 10"
echo " docker run <image> 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
50 changes: 47 additions & 3 deletions src/python_grpc/apps/collector/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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(
Expand All @@ -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()

Expand Down Expand Up @@ -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.")
Expand All @@ -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"),
Expand All @@ -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.")

Expand Down
Loading
Loading