import os import sys import json import uuid import base64 import asyncio import fcntl import atexit 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 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"] # 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'} @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 # 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=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({ "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") 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") webex_room_id = None room_title = f"Support - {client_name}" room_link = None webex_status = "skipped" sms_sent = False 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 # --------------------------------------------------------- # 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/