Files
Smartwave/cloud/app.py
T
Ninluc fed99772fc
Build, push image, and notify Watchtower / build-image (push) Successful in 46s
Build, push image, and notify Watchtower / notify (push) Successful in 8s
Don't override microwave ID
2026-08-24 15:44:22 +02:00

667 lines
24 KiB
Python

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"<p>{gen}</p>"
@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/<job_id>', methods=['POST', 'GET'])
def trigger_job_manually(job_id):
"""Manually triggers any scheduled job immediately by its ID."""
job = scheduler.get_job(job_id)
if not job:
return jsonify({
"error": f"Job '{job_id}' not found",
"available_jobs": [j.id for j in scheduler.get_jobs()]
}), 404
try:
job.func(*job.args)
return jsonify({
"status": "success",
"message": f"Job '{job_id}' executed successfully."
}), 200
except Exception as e:
return jsonify({
"status": "error",
"message": f"Job execution failed: {str(e)}"
}), 500
@app.route("/debug/dishsafety/not_safe", methods=["GET"])
def debug_unsafe_dish():
"""Debug endpoint to simulate an unsafe dish scenario."""
# Simulated unsafe dish data
unsafe_dish = "storage/dishPhotos/dish_78827473286c4f67855f3fbd0ccb9bb4.jpg"
# Call the cooking_params endpoint logic directly
return safety_checker.check_dish_safety(unsafe_dish)
@app.route("/debug/dishsafety/safe", methods=["GET"])
def debug_safe_dish():
"""Debug endpoint to simulate a safe dish scenario."""
# Simulated safe dish data
unsafe_dish = "storage/dishPhotos/dish_af60c26a08a74ce894d453394f417c22.jpg"
# Call the cooking_params endpoint logic directly
return safety_checker.check_dish_safety(unsafe_dish)
# ---------------------------------------------------------
# OData Metadata Definition
# ---------------------------------------------------------
@app.route("/$metadata", methods=["GET"])
def odata_metadata():
"""Provides mandatory OData schema definition for validation."""
xml_metadata = """<?xml version="1.0" encoding="utf-8"?>
<edmx:Edmx Version="4.0" xmlns:edmx="http://docs.oasis-open.org/odata/ns/edmx">
<edmx:DataServices>
<Schema Namespace="MicrowaveNetwork" xmlns="http://docs.oasis-open.org/odata/ns/edm">
<EntityType Name="Telemetry">
<Key><PropertyRef Name="id"/></Key>
<Property Name="id" Type="Edm.String" Nullable="false"/>
<Property Name="device_id" Type="Edm.String"/>
<Property Name="received_at" Type="Edm.DateTimeOffset"/>
</EntityType>
<EntityType Name="CookingParams">
<Key><PropertyRef Name="id"/></Key>
<Property Name="id" Type="Edm.String" Nullable="false"/>
<Property Name="microwave_id" Type="Edm.String"/>
</EntityType>
<EntityType Name="CveAudit">
<Key><PropertyRef Name="id"/></Key>
<Property Name="id" Type="Edm.String" Nullable="false"/>
<Property Name="status" Type="Edm.String"/>
</EntityType>
<EntityType Name="Alert">
<Key><PropertyRef Name="id"/></Key>
<Property Name="id" Type="Edm.String" Nullable="false"/>
<Property Name="orchestrator_id" Type="Edm.String"/>
<Property Name="webex_room_id" Type="Edm.String"/>
<Property Name="webex_status" Type="Edm.String"/>
<Property Name="sms_sent" Type="Edm.Boolean"/>
</EntityType>
<EntityType Name="RfidCard">
<Key><PropertyRef Name="id"/></Key>
<Property Name="id" Type="Edm.String" Nullable="false"/>
<Property Name="card_id" Type="Edm.String" Nullable="false"/>
<Property Name="valid" Type="Edm.Boolean" Nullable="false"/>
</EntityType>
<EntityContainer Name="Container">
<EntitySet Name="Telemetry" EntityType="MicrowaveNetwork.Telemetry"/>
<EntitySet Name="CookingParams" EntityType="MicrowaveNetwork.CookingParams"/>
<EntitySet Name="CveAudit" EntityType="MicrowaveNetwork.CveAudit"/>
<EntitySet Name="Alerts" EntityType="MicrowaveNetwork.Alert"/>
<EntitySet Name="RfidCards" EntityType="MicrowaveNetwork.RfidCard"/>
</EntityContainer>
</Schema>
</edmx:DataServices>
</edmx:Edmx>"""
return Response(xml_metadata, mimetype="application/xml")
if __name__ == "__main__":
app.run(debug=getattr(config, "DEBUG", True))