import os import sys import json import uuid import base64 import asyncio import fcntl import atexit import re from datetime import datetime, timedelta, timezone from flask import Flask, Response, request, jsonify, current_app, abort, g from pymongo import MongoClient, ASCENDING, DESCENDING from apscheduler.schedulers.background import BackgroundScheduler from APIs import generate from APIs.mqtt import send_command from APIs.webex import WebexManager from APIs.shodan import ShodanAuditor from APIs.twilio import send_alert_sms from image_enhancer import enhance_image_for_ai from microwaveCookPlanner import MicrowaveCookPlanner import safety_checker from jobs import job_server_cve_audit, job_client_ip_audit, job_request_telemetry sys.path.insert(0, '..') try: from shared import config except ImportError: from ..shared import config app = Flask(__name__) # --------------------------------------------------------- # Configuration & Setup # --------------------------------------------------------- # Configure MongoDB connection MONGO_URI = os.getenv("MONGO_URI", "mongodb://localhost:27017/") client = MongoClient(MONGO_URI) db = client["microwave_network_db"] # Database Collections client_collection = db["client_data"] cooking_collection = db["cooking_parameters"] telemetry_collection = db["telemetry_data"] alert_collection = db["alert_data"] webex_tokens_collection = db["webex_tokens"] device_ip_collection = db["device_ips"] cve_audit_results = db["cve_audit_results"] rfid_cards_collection = db["rfid_cards"] # Initialize Database Indexes def init_db_indexes(): device_ip_collection.create_index([("received_at", DESCENDING)]) device_ip_collection.create_index([("client_ip", ASCENDING)]) alert_collection.create_index([("orchestrator_id", ASCENDING), ("created_at", DESCENDING)]) client_collection.create_index([("devices.device_id", ASCENDING)]) init_db_indexes() # Environment Configurations WEBEX_CLIENT_ID = os.getenv("WEBEX_CLIENT_ID", "YOUR_WEBEX_CLIENT_ID") WEBEX_CLIENT_SECRET = os.getenv("WEBEX_CLIENT_SECRET", "YOUR_WEBEX_CLIENT_SECRET") WEBEX_REDIRECT_URI = os.getenv("WEBEX_REDIRECT_URI", "https://smartwave.matthiasg.dev/oauth/callback") WEBEX_TEAM_ID = os.getenv("WEBEX_TEAM_ID", "YOUR_WEBEX_TEAM_ID") WEBEX_NINLUC_ID = os.getenv("WEBEX_NINLUC_ID", "YOUR_WEBEX_NINLUC_ID") SAVEUP_TWILIO_API_TOKEN = os.getenv("SAVEUP_TWILIO_API_TOKEN", "false").lower() == "true" # Ensure the camera image storage directory exists when the app starts CAMERA_IMAGE_DIR = "storage/dishPhotos" os.makedirs(CAMERA_IMAGE_DIR, exist_ok=True) # --------------------------------------------------------- # Integrations Setup # --------------------------------------------------------- microwave_cook_planner = MicrowaveCookPlanner() webex_manager = WebexManager( db_collection=webex_tokens_collection, client_id=WEBEX_CLIENT_ID, client_secret=WEBEX_CLIENT_SECRET, redirect_uri=WEBEX_REDIRECT_URI, team_id=WEBEX_TEAM_ID, user_id=WEBEX_NINLUC_ID ) shodan_auditor = ShodanAuditor() # --------------------------------------------------------- # Cron Scheduler (Multi-Worker Safe) # --------------------------------------------------------- scheduler = BackgroundScheduler(daemon=True) def start_scheduler_once(): """Ensures only ONE Gunicorn worker process runs the cron scheduler.""" lock_file_path = "/tmp/scheduler.lock" # Open or create a lock file lock_file = open(lock_file_path, "wb") try: # Request a non-blocking exclusive lock fcntl.flock(lock_file, fcntl.LOCK_EX | fcntl.LOCK_NB) # Register job schedules scheduler.add_job( func=job_server_cve_audit, args=[shodan_auditor, db], trigger="cron", hour=0, minute=0, id="server_cve_audit_job", replace_existing=True ) scheduler.add_job( func=job_client_ip_audit, args=[shodan_auditor, db], trigger="cron", hour="*/6", minute=15, id="client_ip_audit_job", replace_existing=True ) scheduler.add_job( func=job_request_telemetry, args=[send_command], trigger="cron", hour="1,13", minute=0, id="request_telemetry_job", replace_existing=True ) scheduler.start() print(f"[Cron] Scheduler started successfully in Worker PID {os.getpid()}") # Ensure lock release on shutdown def cleanup(): try: fcntl.flock(lock_file, fcntl.LOCK_UN) lock_file.close() except Exception: pass atexit.register(cleanup) except (IOError, OSError): # Lock acquired by another worker - skip starting scheduler print(f"[Cron] Worker PID {os.getpid()} skipped scheduler (already running in another worker).") # Initialize single-instance scheduler start_scheduler_once() # --------------------------------------------------------- # Authentication Middleware # --------------------------------------------------------- EXEMPT_ROUTES = {'hello_world', 'oauth_callback', 'odata_metadata', 'debug_telemetryrequest', 'trigger_job_manually', 'debug_unsafe_dish', 'debug_safe_dish'} @app.before_request def authenticate_request(): if request.endpoint in EXEMPT_ROUTES or request.method == 'OPTIONS': return # Extract API Key from Header (X-API-Key or Bearer token) api_key = request.headers.get("X-API-Key") if not api_key: auth_header = request.headers.get("Authorization", "") if auth_header.startswith("Bearer "): api_key = auth_header.split(" ")[1] if not api_key: return jsonify({ "@odata.error": { "code": "401", "message": "Unauthorized: Missing API Key header" } }), 401 # Validate key against 'devices_authentication' collection device_auth = db.devices_authentication.find_one({"api_key": api_key}) if not device_auth: return jsonify({ "@odata.error": { "code": "401", "message": "Unauthorized: Invalid API Key" } }), 401 # Attach device context (works for both Microwaves and Node-RED) g.device_id = device_auth.get("device_id") # --------------------------------------------------------- # Application Routes # --------------------------------------------------------- @app.route("/") def hello_world(): gen = generate(prompt="Say Hello, to the user !") return f"
{gen}
" @app.route("/oauth/callback") def oauth_callback(): code = request.args.get("code") if not code: return jsonify({"error": "Missing code parameter"}), 400 try: webex_manager.exchange_code(code) return jsonify({"status": "success", "message": "Webex tokens stored successfully!"}), 200 except Exception as e: current_app.logger.exception("Failed to exchange OAuth code") return jsonify({"error": "OAuth exchange failed"}), 500 @app.route("/cooking-params", methods=["POST"]) async def cooking_params(): try: data = request.get_json() if not data or not isinstance(data, dict): return jsonify({"@odata.error": {"code": "400", "message": "Invalid payload"}}), 400 data["microwave_id"] = g.device_id forwarded_for = request.headers.get('X-Forwarded-For') client_ip = forwarded_for.split(',')[0].strip() if forwarded_for else request.remote_addr device_ip_collection.insert_one({ "microwave_id": g.device_id, "client_ip": client_ip, "received_at": datetime.now(timezone.utc) }) height_cm = float(data.get("dish_height", 4.0)) initial_temp_c = float(data.get("ir_initial_temp", 20.0)) microwave_wattage = int(data.get("microwave_wattage", 900)) defrost_mode = bool(data.get("defrost_mode", False)) camera_image_b64 = data.get("camera_image") if not camera_image_b64: return jsonify({"@odata.error": {"code": "400", "message": "Missing camera_image"}}), 400 filename = f"dish_{uuid.uuid4().hex}.jpg" filepath = os.path.join(CAMERA_IMAGE_DIR, filename) with open(filepath, "wb") as f: f.write(base64.b64decode(camera_image_b64)) data["camera_image"] = filepath enhanced_filepath = enhance_image_for_ai(filepath) # Concurrent AI Execution safety_task = asyncio.to_thread(safety_checker.check_dish_safety, filepath) planner_task = asyncio.to_thread( microwave_cook_planner.generate_plan, image_path=enhanced_filepath, height_cm=height_cm, initial_temp_c=initial_temp_c, microwave_wattage=microwave_wattage, defrost_mode=defrost_mode ) safety_result, cook_plan = await asyncio.gather(safety_task, planner_task) data["safety_check"] = safety_result data["analysis_results"] = cook_plan inserted = cooking_collection.insert_one(data) doc_id = str(inserted.inserted_id) # Base OData v4 response wrapper odata_response = { "@odata.context": f"{request.host_url.rstrip('/')}/$metadata#CookingParams/$entity", "@odata.id": f"{request.host_url.rstrip('/')}/cooking-params('{doc_id}')", "id": doc_id, "plan": cook_plan } if not safety_result.get("is_safe", True): odata_response.update({ "error": "Unsafe cooking area", "is_safe": False, "warning_message": safety_result.get("warning_message", "Unsafe materials detected."), "detected_hazards": safety_result.get("detected_hazards", []) }) return jsonify(odata_response), 200 return jsonify(odata_response), 201 except Exception as e: current_app.logger.exception("Error in /cooking-params") return jsonify({"@odata.error": {"code": "500", "message": "Internal processing error"}}), 500 @app.route("/alert", methods=["POST"]) def alert(): data = request.get_json() if data is None or not isinstance(data, dict): return jsonify({ "@odata.error": { "code": "400", "message": "Invalid or missing JSON payload" } }), 400 # Auto-populate or override orchestrator_id using authenticated device context orchestrator_id = data.get("orchestrator_id") or getattr(g, "device_id", None) data["orchestrator_id"] = orchestrator_id data["created_at"] = datetime.now(timezone.utc) try: saved_alert = alert_collection.insert_one(data) doc_id = str(saved_alert.inserted_id) except Exception as e: current_app.logger.exception("Failed to insert alert into database") return jsonify({ "@odata.error": { "code": "500", "message": "Database write failed" } }), 500 alert_info = data.get("alert", {}) alert_type = alert_info.get("type", "Unknown") alert_message = alert_info.get("message", "No message provided") webex_room_id = None room_link = None room_title = None webex_status = "skipped" sms_sent = False if alert_type != "COOKING_SAFETY": client_name = f"Unknown Client (Orchestrator {orchestrator_id})" if orchestrator_id else "Unknown Client" client_email = None client_phone = None if orchestrator_id: client_doc = client_collection.find_one({ "devices": { "$elemMatch": { "device_type": "orchestrator", "device_id": str(orchestrator_id) } } }) if client_doc: client_name = client_doc.get("client_name", client_name) contact_info = client_doc.get("contact_info", {}) client_email = contact_info.get("email") client_phone = contact_info.get("phone") room_title = f"Support - {client_name}" three_minutes_ago = datetime.now(timezone.utc) - timedelta(minutes=3) recent_alert = alert_collection.find_one({ "orchestrator_id": orchestrator_id, "webex_room_id": {"$exists": True, "$ne": None}, "created_at": {"$gte": three_minutes_ago}, "_id": {"$ne": saved_alert.inserted_id} }, sort=[("created_at", DESCENDING)]) if recent_alert: webex_room_id = recent_alert["webex_room_id"] room_title = recent_alert.get("webex_room_title", room_title) try: webex_manager.send_alert_info_message( room_id=webex_room_id, alert_type=alert_type, alert_message=alert_message, alert_id=doc_id, added_alert_info=True ) webex_status = "reused_room" room_link = recent_alert.get("webex_room_link") or f"https://web.webex.com/spaces/{webex_room_id}" except Exception as e: current_app.logger.exception("Failed to send follow-up message to Webex room") webex_status = f"message_failed: {str(e)}" else: try: room_details = webex_manager.create_support_room( client_name, client_email, alert_type, alert_message, doc_id ) if isinstance(room_details, dict): webex_room_id = room_details.get("id") room_title = room_details.get("title", room_title) room_link = room_details.get("meetingLink") or f"https://web.webex.com/spaces/{webex_room_id}" else: webex_room_id = room_details room_link = f"https://web.webex.com/spaces/{webex_room_id}" webex_status = "created_new_room" except Exception as e: current_app.logger.exception("Failed to create Webex support room") webex_status = f"creation_failed: {str(e)}" if webex_room_id and client_phone and not SAVEUP_TWILIO_API_TOKEN: sms_sent = send_alert_sms( to_phone=client_phone, room_title=room_title, room_link=room_link, client_email=client_email ) update_fields = {"sms_sent": sms_sent} if webex_room_id: update_fields.update({ "webex_room_id": webex_room_id, "webex_room_title": room_title, "webex_room_link": room_link }) alert_collection.update_one({"_id": saved_alert.inserted_id}, {"$set": update_fields}) # OData v4 Created Response Payload return jsonify({ "@odata.context": f"{request.host_url.rstrip('/')}/$metadata#Alerts/$entity", "@odata.id": f"{request.host_url.rstrip('/')}/alert('{doc_id}')", "id": doc_id, "status": "success", "message": "Alert saved", "orchestrator_id": orchestrator_id, "webex_room_id": webex_room_id, "webex_room_url": room_link, "webex_status": webex_status, "sms_sent": sms_sent }), 201 @app.route("/telemetry", methods=["POST"]) def telemetry(): data = request.get_json() if data is None: return jsonify({"@odata.error": {"code": "400", "message": "Invalid JSON payload"}}), 400 if isinstance(data, str): try: data = json.loads(data) except (json.JSONDecodeError, TypeError): return jsonify({"@odata.error": {"code": "400", "message": "String payload invalid"}}), 400 data["device_id"] = g.device_id data["received_at"] = datetime.now(timezone.utc).isoformat() try: inserted = telemetry_collection.insert_one(data) doc_id = str(inserted.inserted_id) # OData v4 Created Response Envelope return jsonify({ "@odata.context": f"{request.host_url.rstrip('/')}/$metadata#Telemetry/$entity", "@odata.id": f"{request.host_url.rstrip('/')}/telemetry('{doc_id}')", "id": doc_id, "status": "success", "message": "Telemetry saved", "device_id": g.device_id, "received_at": data["received_at"] }), 201 except Exception as e: current_app.logger.exception("Database write failed") return jsonify({"@odata.error": {"code": "500", "message": "Database write failed"}}), 500 @app.route('/api/security/cve-audit', methods=['POST', 'GET']) def run_cve_audit(): """Protected endpoint triggered by Node-RED using REST API + API Key.""" # Only node-red is allowed in this endpoint if g.device_id != "node-red": return jsonify({ "@odata.error": { "code": "403", "message": "Forbidden: Not allowed to trigger CVE audit" } }), 403 audit_results = shodan_auditor.audit_server_vulnerabilities() inserted = cve_audit_results.insert_one({ "triggered_by": g.device_id, "audit_results": audit_results, "timestamp": datetime.now(timezone.utc) }) doc_id = str(inserted.inserted_id) status_code = 200 if audit_results.get("status") == "PASS" else 409 return jsonify({ "@odata.context": f"{request.host_url.rstrip('/')}/$metadata#CveAudit/$entity", "@odata.id": f"{request.host_url.rstrip('/')}/api/security/cve-audit('{doc_id}')", "id": doc_id, "status": audit_results.get("status"), "audit_results": audit_results, "timestamp": datetime.now(timezone.utc).isoformat() }), status_code @app.route("/rfid-card", methods=["GET"]) def get_rfid_cards(): """OData compliant endpoint for technician RFID card validation.""" filter_str = request.args.get("$filter", "") card_id = request.args.get("card_id") mongo_query = {} # Extract card_id from OData $filter string if not supplied directly if not card_id and filter_str: match = re.search(r"(?:card_id|cardId)\s+eq\s+['\"]([^'\"]+)['\"]", filter_str, re.IGNORECASE) if match: card_id = match.group(1) else: # Generic fallback to extract quoted strings generic_match = re.search(r"['\"]([a-zA-Z0-9]+)['\"]", filter_str) if generic_match: card_id = generic_match.group(1) if card_id: mongo_query["card_id"] = card_id # Filter for active/valid cards only mongo_query["valid"] = True try: cards = list(rfid_cards_collection.find(mongo_query)) result_value = [ { "id": str(card["_id"]), "card_id": card.get("card_id"), "valid": card.get("valid", False) } for card in cards ] return jsonify({ "@odata.context": f"{request.host_url.rstrip('/')}/$metadata#RfidCards", "value": result_value }), 200 except Exception as e: current_app.logger.exception("Error querying rfid_cards collection") return jsonify({"@odata.error": {"code": "500", "message": "Database query error"}}), 500 # --------------------------------------------------------- # Debug Endpoints (Protected) # --------------------------------------------------------- @app.route("/debug/telemetryrequest", methods=["GET"]) def debug_telemetryrequest(): # Construct external HTTP endpoint dynamically based on incoming request host telemetry_url = f"{request.host_url.rstrip('/').replace('http://', 'https://')}/telemetry" cmd_payload = { "action": "request_telemetry", "endpoint": telemetry_url } try: send_command(topic="cmd/all", payload=cmd_payload) return jsonify({ "status": "Telemetry command sent to cmd/all", "published_payload": cmd_payload }), 200 except Exception as e: return jsonify({"error": f"Failed to publish MQTT command: {str(e)}"}), 500 @app.route('/debug/run-job/