Files
Smartwave/orchestrateur/main.py
T
Ninluc 7299a50198
Build, push image, and notify Watchtower / build-image (push) Successful in 1m21s
Build, push image, and notify Watchtower / notify (push) Successful in 17s
Beginning of cooking cycle
2026-07-31 21:50:42 +02:00

395 lines
16 KiB
Python

import base64
import json
import threading
import queue
import time
import traceback
import requests
from orchestrateur.sensors import gps
from shared import get_lora, get_mqtt_client, deviceTypes, config, payloads
from shared.logging import log
from sensors import ultrasonicRanger, temp_hum, button, camera
# --- Read Unique Device ID ---
try:
with open("device_id.txt", "r") as f:
DEVICE_ID = f.read().strip()
except Exception:
try:
with open("/home/pi/SmartWave/orchestrateur/device_id.txt", "r") as f:
DEVICE_ID = f.read().strip()
except Exception:
DEVICE_ID = "RPI_Orchestrateur_Default"
# Thread-safe queue for application messages
data_queue = queue.Queue()
cooking_queue = {}
active_cooks = {}
active_cooks_lock = threading.Lock()
lora = get_lora()
lora.configure()
def lora_listener():
"""Background Thread: Listens to LoRa traffic and responds to Heartbeats."""
print("Thread Écouteur LoRa démarré.")
while True:
paquet = lora.receive_packet(timeout_ms=1000)
if paquet:
donnees = paquet["data"]
expediteur_type = donnees.get("type")
if expediteur_type == deviceTypes.DEVICE_TYPES["MICROWAVE"]:
print(f"\n[Thread LoRa] Heartbeat reçu de {donnees.get('id')}")
reponse = {
"id": DEVICE_ID,
"type": deviceTypes.DEVICE_TYPES["ORCHESTRATOR"]
}
lora.send(reponse)
else:
data_queue.put({"source": "LoRa", "data": paquet})
# --- Setup & Connect MQTT ---
mqtt_client = get_mqtt_client(
host="192.168.50.1", # Using explicit gateway IP to dodge Docker loopback blocks
client_id="smartwave-orchestrateur-"+DEVICE_ID,
use_tls=config.USE_TLS,
cafile="/home/pi/SmartWave/orchestrateur/mqtt/certs/ca.crt",
keepalive=config.MQTT_KEEPALIVE,
)
mqtt_client.connect()
mqtt_client.subscribe(config.MQTT_TOPIC_SENSOR, qos=config.MQTT_QOS)
print(f"Subscribed to topic: {config.MQTT_TOPIC_SENSOR}")
mqtt_client.subscribe(config.MQTT_TOPIC_HELLO, qos=config.MQTT_QOS)
print(f"Subscribed to topic: {config.MQTT_TOPIC_HELLO}")
# --- THE CRUCIAL PAHO FIX ---
# Start Paho's internal background thread. This handles all network packets,
# automatic keepalive pings, and delivery receipts cleanly.
if hasattr(mqtt_client._client, "loop_start"):
mqtt_client._client.loop_start()
print("Paho MQTT asynchronous network loop started.")
def mqtt_listener():
"""Background Thread: Constantly inspects incoming MQTT message cache."""
print("Thread MQTT démarré.")
while True:
message = mqtt_client.get_message()
if message:
# Try to parse the payload as a python dictionary, but if it fails, just print the raw payload
try:
payload = json.loads(message['payload'])
except Exception as e:
print(f"Error parsing MQTT payload: {e}")
payload = message['payload'] # Fallback to raw payload if parsing fails
print(f"\n[Thread MQTT] Message reçu : {message}")
data_queue.put({"source": "MQTT", "topic": message['topic'] ,"data": payload})
# Sleep for 100ms. Prevents the thread from turning into an infinite 100% CPU hog.
time.sleep(0.2)
# Button
button_state = False
def button_callback():
global button_state
button_state = not button_state
print(f"\n[Thread Button] Button state changed to: {button_state}")
button.set_callback(button_callback)
# Launch background monitoring workers
threading.Thread(target=lora_listener, daemon=True).start()
threading.Thread(target=mqtt_listener, daemon=True).start()
# Launch button monitoring thread
button.start_button_monitoring_thread()
print("Orchestrateur prêt. Le main loop est libre.")
def _tryReadSensorsWithRetries(func, exception=True, max_retries=3, delay=1):
"""
Tries to read the sensor max_retries times until the return value of func is not None.
It will then return the value of func. If it fails max_retries times, it will fail if exception is True, otherwise it will return None.
"""
for attempt in range(max_retries):
result = func()
if result is not None:
return result
else:
log(f"Attempt {attempt + 1} failed. Retrying in {delay} seconds...")
time.sleep(delay)
if exception:
raise Exception(f"Failed to read sensor after {max_retries} attempts.")
else:
return None
def read_sensors_for_cooking(microwave_id):
"""Read all sensors and return a dictionary of their values, including the microwave ID."""
log("\nLecture des capteurs...")
sensor_data = {}
sensor_data["microwave_id"] = microwave_id
# === Notify the microwave of needed sensor readings ===
mqtt_client.publish(config.MQTT_TOPIC_COOKING, payloads.mqtt_cooking_init(microwave_id), qos=config.MQTT_QOS)
# Read Ultrasonic Ranger
sensor_data["ultrasonic_distance"] = _tryReadSensorsWithRetries(ultrasonicRanger.get_dish_height)
log(f"\nLecture du capteur Ultrason : {sensor_data['ultrasonic_distance']}")
# Read Temperature and Humidity
temperature, humidity = temp_hum.get_temperature_and_humidity_with_retry()
if temperature is not None and humidity is not None:
log(f"\nLecture du capteur Temp/Hum : {temperature}, {humidity}")
sensor_data["temperature"] = temperature
sensor_data["humidity"] = humidity
# Camera
def _getPicture():
picture_bytes = None
try:
picture_bytes = camera.get_picture()
return picture_bytes
except Exception as e:
log(f"Error reading camera data: {e}")
return None
sensor_data["camera_image"] = _tryReadSensorsWithRetries(_getPicture, exception=True)
log(f"\nPhoto de la Caméra : {len(sensor_data['camera_image'])} bytes")
# Read Button State (last because he can still change state while reading other sensors)
sensor_data["defrost_mode"] = button_state
cooking_queue[microwave_id] = sensor_data
def _stop_hardware(microwave_id: str):
"""
Hardware driver stop — halts magnetron/turntable immediately.
"""
print(f"[{microwave_id}] 🛑 Emergency stop issued to hardware.")
# TODO: Add physical hardware stop command here
# e.g., gpio_controller.stop()
def _send_params_to_microwave(microwave_id: str, cook_time: int, power_level: int, target_temp: float, cancel_event: threading.Event):
"""
Triggers physical microwave execution.
"""
if cancel_event.is_set():
return
print(f"[{microwave_id}] Sending microwave {microwave_id} cooking parameters : {cook_time}s @ {power_level}W power, target temp {target_temp}°C.")
mqtt_client.publish(config.MQTT_TOPIC_COOKING, payloads.mqtt_cooking_config(microwave_id, cook_time, power_level, target_temp), qos=config.MQTT_QOS)
def _cooking_worker(microwave_id: str, sensors_data: dict, cancel_event: threading.Event):
"""Worker function executing cloud API calls and hardware triggers."""
URL = "https://smartwave.matthiasg.dev/cooking-params"
try:
# Check cancellation before network call
if cancel_event.is_set():
print(f"[{microwave_id}] Job canceled before starting API call.")
return
log(f"[{microwave_id}] Sending sensor data to cloud API...")
# Ensure camera_image is encoded to Base64 string if it's currently raw bytes
if isinstance(sensors_data.get("camera_image"), bytes):
sensors_data["camera_image"] = base64.b64encode(sensors_data["camera_image"]).decode("utf-8")
# 1. HTTP Request (15-second timeout)
response = requests.post(URL, json=sensors_data, timeout=360) # timeout for long-running requests
# Check cancellation right after network call returns
if cancel_event.is_set():
print(f"[{microwave_id}] Job was canceled while waiting for cloud response. Discarding result.")
return
log(f"[{microwave_id}] Cloud API responded with : {response.json()}")
if response.status_code != 200 and response.status_code != 201:
print(f"[{microwave_id}] Cloud API returned error {response.status_code}: {response.json()}")
response.raise_for_status()
# 2. Extract Response Parameters
response_json = response.json()
cook_plan = response_json.get("cook_plan", {})
cook_time = cook_plan.get("cook_time_seconds")
power_level = cook_plan.get("effective_power_watts")
target_temp = cook_plan.get("target_temp")
dish_name = response_json.get("dish_name", "Unknown Dish")
if cook_time is None or power_level is None or target_temp is None:
print(f"[{microwave_id}] Cloud returned incomplete plan: {response_json}")
return
# Check cancellation before starting physical microwave
if cancel_event.is_set():
print(f"[{microwave_id}] Job was canceled before starting hardware execution.")
return
print(f"[{microwave_id}] Received plan for '{dish_name}': {cook_time}s @ {power_level}W power, target temp {target_temp}°C.")
# 3. Start Hardware Execution
_send_params_to_microwave(microwave_id, cook_time, power_level, target_temp, cancel_event)
except requests.exceptions.Timeout:
print(f"[{microwave_id}] Request timed out waiting for cloud response.")
except requests.exceptions.RequestException as e:
print(f"[{microwave_id}] HTTP error reaching cloud API: {e}")
except Exception as e:
print(f"[{microwave_id}] Unexpected error in worker thread: {e}")
traceback.print_exc()
finally:
# Clean up registry entry if this worker was the active one
with active_cooks_lock:
if active_cooks.get(microwave_id) == cancel_event:
del active_cooks[microwave_id]
def start_cooking_for_microwave(microwave_id: str, sensors_data: dict):
"""
Sends sensors data to the cloud and starts cooking in a separate thread.
If a worker is already running for the given microwave_id, it cancels
the previous process and stops the hardware before starting the new one.
"""
with active_cooks_lock:
# 1. If an active job exists for this microwave, cancel it
if microwave_id in active_cooks:
print(f"[{microwave_id}] Existing cooking job detected! Canceling old worker...")
active_cooks[microwave_id].set() # Signal existing thread to abort
_stop_hardware(microwave_id) # Stop hardware immediately
# 2. Register a new cancellation event for this microwave
cancel_event = threading.Event()
active_cooks[microwave_id] = cancel_event
# 3. Start the new background worker thread
thread = threading.Thread(
target=_cooking_worker,
args=(microwave_id, sensors_data, cancel_event),
daemon=True
)
thread.start()
# Sensor reading
def read_sensors():
"""Read all sensors and return a dictionary of their values."""
log("\nLecture des capteurs...")
sensor_data = {}
# Read Ultrasonic Ranger
distance = ultrasonicRanger.get_dish_height()
if distance is not None:
log(f"\nLecture du capteur Ultrason : {distance}")
sensor_data["ultrasonic_distance"] = distance
# Read Temperature and Humidity
temperature, humidity = temp_hum.get_temperature_and_humidity()
if temperature is not None and humidity is not None:
log(f"\nLecture du capteur Temp/Hum : {temperature}, {humidity}")
sensor_data["temperature"] = temperature
sensor_data["humidity"] = humidity
# Read GPS Data
gps_data = gps.get_gps_data()
if gps_data:
log(f"\nLecture du capteur GPS : {gps_data}")
sensor_data["gps"] = gps_data
# Camera
picture_bytes = None
try:
picture_bytes = camera.get_picture()
log(f"\nLecture du capteur Caméra : {len(picture_bytes)} bytes")
sensor_data["camera_image"] = picture_bytes
except Exception as e:
log(f"Error reading camera data: {e}")
# Read Button State (last because he can still change state while reading other sensors)
sensor_data["defrost_state"] = button_state
return sensor_data
# --- MAIN EXECUTION LOOP ---
while True:
try:
# === TREAT MESSAGE QUEUE ===
try:
msg = data_queue.get(block=False)
# print(msg)
if msg["source"] == "LoRa":
print(f"\n[Main Loop] LoRa : Données traitées : {msg['data']}")
elif msg["source"] == "MQTT":
# MQTT HELLO
if (msg["topic"] == config.MQTT_TOPIC_HELLO.decode('utf-8')):
if ("id_orchestrator" in msg["data"] and msg["data"]["id_orchestrator"] == DEVICE_ID):
# Do not answer to messages coming from me
continue
microwave_id = msg["data"]["id_microwave"]
print(f"\n[Main Loop] MQTT : Hello reçu de {microwave_id}.")
# Responds
mqtt_client.publish(config.MQTT_TOPIC_HELLO, payloads.mqtt_hello_ack(DEVICE_ID, microwave_id), qos=config.MQTT_QOS)
print(f"[Main Loop] MQTT : Réponse Hello envoyée à {microwave_id}.")
# TODO : Save in database
# MQTT SENSOR DATA
elif (msg["topic"] == config.MQTT_TOPIC_SENSOR.decode('utf-8')):
print(f"\n[Main Loop] MQTT : Données capteurs reçues du micro-ondes : {msg['data']}")
microwave_id = msg["data"].get("id_microwave")
if not microwave_id:
print("[Main Loop] MQTT : Données capteurs reçues sans ID micro-ondes. Ignoré.")
continue
# Get the already existing cooking data for this microwave
sensors_data = cooking_queue.get(microwave_id)
if sensors_data is None:
print(f"[Main Loop] MQTT : Données capteurs reçues pour {microwave_id} mais aucune donnée de cuisson en cours. Ignoré.")
continue
# Merge the received sensor data into the existing cooking data
sensors_data["ir_initial_temp"] = msg["data"].get("dish_temp")
sensors_data["ir_ambient_temp"] = msg["data"].get("ambient_temp")
start_cooking_for_microwave(microwave_id, sensors_data)
# Remove the cooking data from the queue since it's now being processed
del cooking_queue[microwave_id]
print(f"\n[Main Loop] MQTT : Données traitées : {msg['data']}")
except queue.Empty:
pass
# === CHECK FOR DISH INSERTED ===
# Read the dish height from the ultrasonic sensor. If it's below a certain threshold, we assume a dish has been inserted.
dish_height = ultrasonicRanger.get_dish_height()
if dish_height is not None and dish_height > 2.0: # Threshold in cm for detecting a dish
print(f"\n[Main Loop] Dish detected at height: {dish_height} cm. Initiating sensor read...")
# Read all sensors and store the data in the cooking queue for this microwave
read_sensors_for_cooking("2")
print(f"[Main Loop] Sensor data collected and queued for cooking.")
# DEBUG : Read sensors
# sensor_values = read_sensors()
# if sensor_values:
# sensor_values_print = sensor_values.copy()
# if "camera_image" in sensor_values_print:
# sensor_values_print["camera_image"] = f"<{len(sensor_values_print['camera_image'])} bytes>"
# print(f"\nCapteurs Données lues : {sensor_values_print}")
time.sleep(3)
except KeyboardInterrupt:
break
except Exception as e:
traceback.print_exc()
time.sleep(1) # Prevents rapid error logging in case of persistent issues
# Clean termination
if hasattr(mqtt_client._client, "loop_stop"):
mqtt_client._client.loop_stop()
mqtt_client.close()