commit 57c9e8dbbf58727221cd9faff0d27f1cc40df9b8 Author: Suriya Date: Mon Jul 27 19:55:10 2026 +0530 Add MQTT webhook bridge with resilient MQTT connection handling FastAPI service that subscribes to nearle/# on the MQTT broker and forwards rider status, periodic logs, and profile updates to the Jupiter webhook API. Includes fixes for a failure mode that took the bridge down silently for ~3.8 days in June 2026, during which it served HTTP 200 on / while forwarding zero events: - Retry MQTT connection with backoff instead of letting a single failed connect kill the background thread permanently. - Subscribe inside on_connect so subscriptions are restored after an automatic reconnect; the session is clean, so a reconnect previously came back subscribed to nothing. - Add /healthz returning 503 when MQTT is disconnected, so a monitor can distinguish a live bridge from a dead one. - Replace thread-per-message forwarding with a bounded worker pool, which keeps memory flat when the upstream API stalls on its 5s timeout. - Pin requirements.txt to the versions verified in production; it was unpinned and had already drifted to paho-mqtt 2.1.0. - Set PYTHONUNBUFFERED so container logs are not block-buffered. - Name the compose image nearle-api-image to match the deployed tag. Co-Authored-By: Claude Opus 5 (1M context) diff --git a/.gitattributes b/.gitattributes new file mode 100644 index 0000000..ed8737b --- /dev/null +++ b/.gitattributes @@ -0,0 +1,3 @@ +# This code is built and run in a Linux container; keep LF everywhere so a +# Windows checkout does not push CRLF back to the server. +* text=auto eol=lf diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..88eb0bf --- /dev/null +++ b/.gitignore @@ -0,0 +1,4 @@ +__pycache__/ +*.py[cod] +.env +*.bak diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..23df470 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,16 @@ +FROM python:3.9-slim + +WORKDIR /app + +# Without this, stdout is block-buffered when not a TTY, so `docker logs` lags far +# behind reality — during the June outage the "MQTT Connection failed" line did not +# surface until long after the failure. +ENV PYTHONUNBUFFERED=1 + +COPY requirements.txt . + +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"] diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..9547228 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,31 @@ +services: + nearle-api: + build: . + # Pinned so `docker compose up -d --build` rebuilds the SAME tag the running + # container was created from (it was started with `docker run`, not compose). + # Without this, compose would build `nearleapi-nearle-api` and leave a + # confusing second image behind. + image: nearle-api-image + container_name: nearle_api + restart: unless-stopped + ports: + - "8000:8000" + environment: + - MQTT_BROKER=66.116.225.226 + - API_STATUS_URL=https://jupiter.nearle.app/live/api/v1/utils/createriderstatus + - API_LOGS_URL=https://jupiter.nearle.app/live/api/v1/utils/createriderperiodiclog + - API_PROFILE_URL=https://jupiter.nearle.app/live/api/v1/utils/createuserredis + # /healthz returns 503 whenever the MQTT side is down, so `docker ps` shows + # "unhealthy" instead of the container looking fine while forwarding nothing. + healthcheck: + test: ["CMD-SHELL", "python -c \"import urllib.request,sys; sys.exit(0 if urllib.request.urlopen('http://127.0.0.1:8000/healthz', timeout=3).status==200 else 1)\""] + interval: 60s + timeout: 10s + retries: 3 + start_period: 30s + networks: + - nearle-net + +networks: + nearle-net: + external: true diff --git a/main.py b/main.py new file mode 100644 index 0000000..be77f3f --- /dev/null +++ b/main.py @@ -0,0 +1,206 @@ +from fastapi import FastAPI +import paho.mqtt.client as mqtt +from paho.mqtt.enums import CallbackAPIVersion +import json +import threading +import os +import time +import requests +import datetime +from concurrent.futures import ThreadPoolExecutor + +app = FastAPI(title="Nearle MQTT Webhook Bridge API") + +# Configuration +MQTT_BROKER = os.getenv("MQTT_BROKER", "66.116.225.226") +MQTT_PORT = int(os.getenv("MQTT_PORT", "1883")) +MQTT_USER = os.getenv("MQTT_USER", "admin") +MQTT_PASS = os.getenv("MQTT_PASS", "Package@321#") +MQTT_TOPIC = os.getenv("MQTT_TOPIC", "nearle/#") + +# Define your specific external POST APIs here +API_STATUS_URL = os.getenv("API_STATUS_URL", "https://jupiter.nearle.app/live/api/v1/utils/createriderstatus") +API_LOGS_URL = os.getenv("API_LOGS_URL", "https://jupiter.nearle.app/live/api/v1/utils/createriderperiodiclog") +API_PROFILE_URL = os.getenv("API_PROFILE_URL", "https://jupiter.nearle.app/live/api/v1/utils/createuserredis") + +# Bounded worker pool. Previously every MQTT message spawned its own thread, so a +# stall on the Jupiter API (the 5s read timeouts seen in the logs) let threads pile +# up without limit. A fixed pool keeps memory flat and applies natural backpressure. +FORWARD_WORKERS = int(os.getenv("FORWARD_WORKERS", "16")) +_executor = ThreadPoolExecutor(max_workers=FORWARD_WORKERS, thread_name_prefix="forward") + +# Liveness state, exposed on / so a monitor can tell a live bridge from a dead one. +_state = { + "mqtt_connected": False, + "last_connect_at": None, + "last_disconnect_at": None, + "last_message_at": None, + "connect_attempts": 0, + "messages_received": 0, +} +_state_lock = threading.Lock() + + +def _now(): + return datetime.datetime.now(datetime.timezone.utc).isoformat() + + +# Function to forward data to a specific URL +def forward_to_api(target_url: str, payload_str: str): + if not target_url: + return + + # Determine the content type and prepare data + headers = {} + + # Check if the payload is JSON or plain string + try: + payload_data = json.loads(payload_str) + headers['Content-Type'] = 'application/json' + + # We send the exact parsed JSON so the receiving API sees it cleanly + data_kwargs = {"json": payload_data} + except: + # If it's a plain string like "Offline", send as text/plain + headers['Content-Type'] = 'text/plain' + data_kwargs = {"data": payload_str} + + try: + response = requests.post(target_url, headers=headers, timeout=5, **data_kwargs) + if response.status_code in (200, 201): + print(f"SUCCESS ({response.status_code}) -> {target_url} | Sent: {data_kwargs}", flush=True) + else: + print(f"WARNING: HTTP {response.status_code} from {target_url} | Response: {response.text}", flush=True) + except Exception as e: + print(f"ERROR: Failed to send to {target_url} | Exception: {e}", flush=True) + + +def route_payload(data_type: str, rider_id: str, payload_str: str): + payload_str = payload_str.strip() + if data_type == "status": + # Ensure userid is an integer if possible so Jupiter API doesn't reject it + r_id = int(rider_id) if rider_id.isdigit() else rider_id + status_json = json.dumps({"userid": r_id, "status": payload_str}) + forward_to_api(API_STATUS_URL, status_json) + elif data_type == "logs": + forward_to_api(API_LOGS_URL, payload_str) + elif data_type == "profile": + forward_to_api(API_PROFILE_URL, payload_str) + else: + # Ignore other data types or add additional routes here + pass + + +# MQTT Callbacks +def on_connect(client, userdata, connect_flags, reason_code, properties=None): + if reason_code == 0: + with _state_lock: + _state["mqtt_connected"] = True + _state["last_connect_at"] = _now() + # Subscribing here (not once after connect) is the important part: the + # session is clean, so every reconnect starts with no subscriptions. Doing + # it in on_connect means the bridge resubscribes automatically after any + # broker restart or network blip instead of sitting connected and deaf. + client.subscribe(MQTT_TOPIC) + print(f"MQTT connected to {MQTT_BROKER}:{MQTT_PORT}, subscribed to {MQTT_TOPIC}", flush=True) + else: + print(f"MQTT connect refused: {reason_code}", flush=True) + + +def on_disconnect(client, userdata, disconnect_flags, reason_code, properties=None): + with _state_lock: + _state["mqtt_connected"] = False + _state["last_disconnect_at"] = _now() + print(f"MQTT disconnected: {reason_code} (will auto-reconnect)", flush=True) + + +def on_message(client, userdata, msg): + try: + topic = msg.topic + payload = msg.payload.decode() + + with _state_lock: + _state["last_message_at"] = _now() + _state["messages_received"] += 1 + + # Topic structure: nearle/riders/{id}/{data_type}/... + parts = topic.split('/') + if len(parts) >= 4 and parts[0] == "nearle" and parts[1] == "riders": + rider_id = parts[2] + data_type = parts[3] + + # Hand off to the pool so the MQTT network loop is never blocked + _executor.submit(route_payload, data_type, rider_id, payload) + + except Exception as e: + print(f"Error processing MQTT message: {e}", flush=True) + + +def start_mqtt(): + # This loop is why the bridge no longer dies silently. The old version called + # connect() inside a try/except: one failure ("No route to host" on 2026-06-10) + # killed the thread for good while uvicorn kept answering 200 OK on /, so the + # bridge looked healthy while forwarding nothing until someone restarted it. + backoff = 5 + while True: + try: + with _state_lock: + _state["connect_attempts"] += 1 + + client = mqtt.Client(CallbackAPIVersion.VERSION2) + client.username_pw_set(MQTT_USER, MQTT_PASS) + client.on_connect = on_connect + client.on_disconnect = on_disconnect + client.on_message = on_message + client.reconnect_delay_set(min_delay=1, max_delay=60) + + print(f"Connecting to MQTT Broker at {MQTT_BROKER}...", flush=True) + + # connect_async + retry_first_connection means a broker that is down at + # startup is retried instead of being a fatal error. + client.connect_async(MQTT_BROKER, MQTT_PORT, keepalive=60) + backoff = 5 + client.loop_forever(retry_first_connection=True) + + print("MQTT loop returned unexpectedly; restarting.", flush=True) + except Exception as e: + print(f"MQTT loop crashed: {e} | retrying in {backoff}s", flush=True) + finally: + with _state_lock: + _state["mqtt_connected"] = False + + time.sleep(backoff) + backoff = min(backoff * 2, 60) + + +# Start MQTT in a background thread +@app.on_event("startup") +async def startup_event(): + thread = threading.Thread(target=start_mqtt, daemon=True) + thread.start() + + +# simple health check endpoint +@app.get("/") +def read_root(): + with _state_lock: + snapshot = dict(_state) + return { + "status": "Nearle Webhook Bridge is running", + "description": "Listening to MQTT and forwarding to external POST APIs.", + "mqtt": snapshot, + } + + +@app.get("/healthz") +def healthz(): + """Returns 200 only when the MQTT side is actually connected, so an external + monitor can detect a bridge that is up but not receiving anything.""" + with _state_lock: + connected = _state["mqtt_connected"] + snapshot = dict(_state) + from fastapi.responses import JSONResponse + return JSONResponse( + status_code=200 if connected else 503, + content={"healthy": connected, "mqtt": snapshot}, + ) diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..ba24648 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,7 @@ +# Pinned to the versions verified running in production on 2026-07-27. +# Previously unpinned, so every rebuild silently pulled whatever was latest. +fastapi==0.128.8 +uvicorn==0.39.0 +paho-mqtt==2.1.0 +requests==2.32.5 +pydantic==2.12.5