Timeout before sending + fix api
This commit is contained in:
@@ -6,7 +6,7 @@ import json
|
|||||||
|
|
||||||
EDAMAM_APP_ID__FOOD = os.getenv("EDAMAM_APP_ID__FOOD", None)
|
EDAMAM_APP_ID__FOOD = os.getenv("EDAMAM_APP_ID__FOOD", None)
|
||||||
EDAMAM_APP_KEY__FOOD = os.getenv("EDAMAM_APP_KEY__FOOD", None)
|
EDAMAM_APP_KEY__FOOD = os.getenv("EDAMAM_APP_KEY__FOOD", None)
|
||||||
SAVE_EDAMAM_API_TOKEN = os.getenv("SAVE_EDAMAM_API_TOKEN", False)
|
SAVEUP_EDAMAM_API_TOKEN = os.getenv("SAVEUP_EDAMAM_API_TOKEN", False)
|
||||||
|
|
||||||
class EdamamAPI:
|
class EdamamAPI:
|
||||||
"""Wrapper around Edamam Food Database API 2.0 (Vision & Nutrients)"""
|
"""Wrapper around Edamam Food Database API 2.0 (Vision & Nutrients)"""
|
||||||
@@ -17,7 +17,7 @@ class EdamamAPI:
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
def analyze_dish_image(self, image_file_path: str):
|
def analyze_dish_image(self, image_file_path: str):
|
||||||
if SAVE_EDAMAM_API_TOKEN:
|
if SAVEUP_EDAMAM_API_TOKEN:
|
||||||
time.sleep(3)
|
time.sleep(3)
|
||||||
return json.loads("""
|
return json.loads("""
|
||||||
{
|
{
|
||||||
|
|||||||
+12
-2
@@ -3,6 +3,7 @@ import base64
|
|||||||
import uuid
|
import uuid
|
||||||
import datetime
|
import datetime
|
||||||
import sys
|
import sys
|
||||||
|
import json
|
||||||
from flask import Flask, request, jsonify
|
from flask import Flask, request, jsonify
|
||||||
from pymongo import MongoClient
|
from pymongo import MongoClient
|
||||||
from APIs import generate, EdamamAPI
|
from APIs import generate, EdamamAPI
|
||||||
@@ -112,14 +113,23 @@ def cooking_params():
|
|||||||
def telemetry():
|
def telemetry():
|
||||||
data = request.get_json()
|
data = request.get_json()
|
||||||
|
|
||||||
if not data:
|
if data is None:
|
||||||
return jsonify({"error": "Invalid or missing JSON payload"}), 400
|
return jsonify({"error": "Invalid or missing JSON payload"}), 400
|
||||||
|
|
||||||
|
# If the payload was double-encoded as a string, deserialize it
|
||||||
|
if isinstance(data, str):
|
||||||
|
try:
|
||||||
|
data = json.loads(data)
|
||||||
|
except (json.JSONDecodeError, TypeError):
|
||||||
|
return jsonify({"error": "String payload could not be parsed as JSON"}), 400
|
||||||
|
|
||||||
|
if not isinstance(data, dict):
|
||||||
|
return jsonify({"error": "Expected a JSON object/dictionary"}), 400
|
||||||
|
|
||||||
# Stamp UTC timestamp for Node-RED queries
|
# Stamp UTC timestamp for Node-RED queries
|
||||||
data["received_at"] = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
data["received_at"] = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# Mongo creates '_id' automatically upon insertion
|
|
||||||
telemetry_collection.insert_one(data)
|
telemetry_collection.insert_one(data)
|
||||||
return jsonify({"status": "success", "message": "Telemetry saved"}), 200
|
return jsonify({"status": "success", "message": "Telemetry saved"}), 200
|
||||||
|
|
||||||
|
|||||||
@@ -415,10 +415,10 @@ async def handle_telemetry_request(endpoint: str):
|
|||||||
# 1. Fetch systemd logs asynchronously
|
# 1. Fetch systemd logs asynchronously
|
||||||
logs_list = await get_systemd_logs(lines=200, service_name="smartwave")
|
logs_list = await get_systemd_logs(lines=200, service_name="smartwave")
|
||||||
|
|
||||||
# 2. Unpack temperature and humidity from sensors.temp_hum
|
# 2. Unpack temperature and humidity
|
||||||
ambient_temp, ambient_humidity = temp_hum.get_temperature_and_humidity_with_retry()
|
ambient_temp, ambient_humidity = temp_hum.get_temperature_and_humidity_with_retry()
|
||||||
|
|
||||||
# 3. Build the telemetry payload
|
# 3. Build the telemetry payload (returns a JSON string)
|
||||||
telemetry_payload = payloads.telemetry_payload(
|
telemetry_payload = payloads.telemetry_payload(
|
||||||
device_id=DEVICE_ID,
|
device_id=DEVICE_ID,
|
||||||
microwave_states=microwave_states,
|
microwave_states=microwave_states,
|
||||||
@@ -431,14 +431,18 @@ async def handle_telemetry_request(endpoint: str):
|
|||||||
logs=logs_list
|
logs=logs_list
|
||||||
)
|
)
|
||||||
|
|
||||||
print(telemetry_payload)
|
|
||||||
print(f"[Telemetry] Sending payload with {len(logs_list)} log entries...")
|
print(f"[Telemetry] Sending payload with {len(logs_list)} log entries...")
|
||||||
|
|
||||||
# 4. Offload blocking HTTP POST to thread pool
|
# 4. Wait for he microwave turn to send the telemetry data to the cloud endpoint
|
||||||
|
timeout = int(DEVICE_ID.split("_")[-1]) * config.TELEMETRY_SEND_INTERVAL
|
||||||
|
print(f"[Telemetry] Waiting for {timeout}s before sending telemetry to avoid collisions...")
|
||||||
|
await asyncio.sleep(timeout)
|
||||||
|
|
||||||
|
# 5. Offload blocking HTTP POST to thread pool
|
||||||
response = await asyncio.to_thread(
|
response = await asyncio.to_thread(
|
||||||
requests.post,
|
requests.post,
|
||||||
endpoint,
|
endpoint,
|
||||||
json=telemetry_payload,
|
data=telemetry_payload, # <-- Changed from json=telemetry_payload
|
||||||
headers={"Content-Type": "application/json"},
|
headers={"Content-Type": "application/json"},
|
||||||
timeout=15
|
timeout=15
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -19,3 +19,6 @@ MQTT_HELLO_INTERVAL = 30
|
|||||||
COOKING_COMPARTMENT_HEIGHT = 30 # cm
|
COOKING_COMPARTMENT_HEIGHT = 30 # cm
|
||||||
|
|
||||||
BUZZER_ACTIVATED = False
|
BUZZER_ACTIVATED = False
|
||||||
|
|
||||||
|
# TELEMETRY
|
||||||
|
TELEMETRY_SEND_INTERVAL = 3
|
||||||
Reference in New Issue
Block a user