diff --git a/src/deployment/docker/docker-compose.yml b/src/deployment/docker/docker-compose.yml index a6815170f..e0de62c83 100644 --- a/src/deployment/docker/docker-compose.yml +++ b/src/deployment/docker/docker-compose.yml @@ -110,7 +110,9 @@ services: environment: MAIL_STARTTLS: "true" MAIL_SSL_TLS: "false" - + MQTT_BROKER_URL: "ts-mqtt-server-cont" + MQTT_BROKER_PORT: "1883" + MQTT_PUBLISH_URL: "projectecho/engine/2" stdin_open: false tty: true diff --git a/src/production/backend/app/main.py b/src/production/backend/app/main.py index 5e81b6163..3699e4d8f 100644 --- a/src/production/backend/app/main.py +++ b/src/production/backend/app/main.py @@ -67,6 +67,34 @@ ) +######### MQTT live event handling (FR-A2) +from app.services.mqtt_client import start_mqtt_client, get_connection_state, get_latest_events + +@app.on_event("startup") +def startup_mqtt(): + start_mqtt_client() + +@app.get("/mqtt/connection-state", tags=["mqtt"]) +def mqtt_connection_state(): + return {"state": get_connection_state()} + +@app.get("/mqtt/latest-events", tags=["mqtt"]) +def mqtt_latest_events(): + from app.routers.hmi import show_latest_events + + events = [ + {"eventType": "vocalization", **event} + for event in show_latest_events(limit=20) + ] + events += [ + event for event in get_latest_events() + if event.get("eventType") != "vocalization" + ] + + return {"events": events} + + + # app.include_router(hmi.router, tags=['hmi'], prefix='/hmi') # app.include_router(engine.router, tags=['engine'], prefix='/engine') # app.include_router(sim.router, tags=['sim'], prefix='/sim') diff --git a/src/production/backend/app/routers/hmi.py b/src/production/backend/app/routers/hmi.py index 32bf7187a..543d7db43 100644 --- a/src/production/backend/app/routers/hmi.py +++ b/src/production/backend/app/routers/hmi.py @@ -134,6 +134,34 @@ def show_event_from_time(start: str, end: str): events = serializers.eventSpeciesListEntity(Events.aggregate(aggregate)) return events +@router.get("/latest_events", response_description="Get most recent classified detection events") +def show_latest_events(limit: int = 20): + aggregate = [ + { "$sort": { "timestamp": -1 } }, + { "$limit": limit }, + { + '$lookup': { + 'from': 'species', + 'localField': 'species', + 'foreignField': '_id', + 'as': 'info' + } + }, + { + "$replaceRoot": { "newRoot": { "$mergeObjects": [ { "$arrayElemAt": [ "$info", 0 ] }, "$$ROOT" ] } } + }, + { + '$project': { "audioClip": 0, "sampleRate": 0} + }, + { + "$addFields": { + "timestamp": { "$toLong": "$timestamp" } + } + } + ] + events = serializers.eventSpeciesListEntity(Events.aggregate(aggregate)) + return events + @router.get("/audio", response_description="returns audio given ID") def show_audio(id: str): aggregate = [ diff --git a/src/production/backend/app/services/mqtt_client.py b/src/production/backend/app/services/mqtt_client.py new file mode 100644 index 000000000..c6d6f2010 --- /dev/null +++ b/src/production/backend/app/services/mqtt_client.py @@ -0,0 +1,133 @@ +import paho.mqtt.client as mqtt +import os +import json + +connection_state = "disconnected" # disconnected | connecting | connected | reconnecting +latest_events = {} + +MQTT_BROKER_URL = os.environ.get("MQTT_BROKER_URL", "ts-mqtt-server-cont") +MQTT_BROKER_PORT = int(os.environ.get("MQTT_BROKER_PORT", 1883)) +MQTT_TOPICS = os.environ.get( + "MQTT_TOPICS", + "projectecho/engine/2,projectecho/movement,iot/data/test", +).split(",") + +def on_connect(client, userdata, flags, rc): + global connection_state + connection_state = "connected" + print(f"[MQTT] Connected, rc={rc}") + for topic in MQTT_TOPICS: + client.subscribe(topic.strip()) + +def on_disconnect(client, userdata, rc): + global connection_state + connection_state = "reconnecting" + print(f"[MQTT] Disconnected, rc={rc} - attempting reconnect") + +def on_message(client, userdata, msg): + normalized = normalize_payload(msg.payload, msg.topic) + if normalized: + print(f"[MQTT] Normalized event: {normalized}") + key = normalized.get("_id", "unknown") + latest_events[key] = normalized + # TODO: forward `normalized` to wherever the frontend/dashboard reads from + +def normalize_payload(raw_payload, topic): + """ + Converts a raw MQTT message into a shape the frontend can use directly. + Routes by event type: vocalization, movement, sensor_health, iot_node. + """ + try: + data = json.loads(raw_payload) + except (json.JSONDecodeError, TypeError): + print(f"[MQTT] Could not parse payload on {topic} as JSON") + return None + + # Vocalization / recording events (from comms_manager.py's + # mqtt_send_random_audio_msg and mqtt_send_recording_msg). + # NOTE: species classification (species/commonName/status/diet) is NOT + # available at this stage — it only exists after the engine classifies + # the audio and posts the result to MongoDB via HTTP, not over MQTT. + # We use clear placeholders here rather than fabricate values. + if "animalEstLLA" in data or "audioClip" in data: + return { + "eventType": "vocalization", + "_id": f"{data.get('sensorId', 'unknown')}_{data.get('timestamp', '')}", + "timestamp": data.get("timestamp"), + "confidence": data.get("animalLLAUncertainty"), + "species": "unclassified", + "commonName": "unclassified", + "type": "mammal", + "status": "normal", + "diet": "unknown", + "animalLLAUncertainty": data.get("animalLLAUncertainty"), + "animalEstLLA": data.get("animalEstLLA"), + "animalTrueLLA": data.get("animalTrueLLA"), + "sensorId": data.get("sensorId"), + "microphoneLLA": data.get("microphoneLLA"), + } + + # Movement events + if "animalId" in data and "species" in data: + return { + "eventType": "movement", + "_id": f"{data.get('animalId', 'unknown')}_{data.get('timestamp', '')}", + "timestamp": data.get("timestamp"), + "animalId": data.get("animalId"), + "species": data.get("species"), + "animalTrueLLA": data.get("animalTrueLLA"), + # Placeholders: real species metadata isn't available on this + # MQTT payload, only after DB enrichment (see vocalization + # branch above for the same limitation). + "type": "mammal", + "status": "normal", + "diet": "omnivore", + } + + # Sensor health events + if "cpu" in data or "batteryPct" in data: + return { + "eventType": "sensor_health", + "_id": f"{data.get('sensorId', 'unknown')}_{data.get('timestamp', '')}", + "timestamp": data.get("timestamp"), + "sensorId": data.get("sensorId"), + "status": data.get("status"), + "batteryPct": data.get("batteryPct"), + "cpu": data.get("cpu"), + "ram": data.get("ram"), + } + + # IoT node updates + if "nodeId" in data: + return { + "eventType": "iot_node", + "_id": f"{data.get('nodeId', 'unknown')}_{data.get('timestamp', '')}", + "timestamp": data.get("timestamp"), + "nodeId": data.get("nodeId"), + "status": data.get("status"), + } + + print(f"[MQTT] Unrecognized payload shape on {topic}: {list(data.keys())}") + return {"eventType": "unknown", "raw": data} + +def start_mqtt_client(): + global connection_state + connection_state = "connecting" + client = mqtt.Client() + client.on_connect = on_connect + client.on_disconnect = on_disconnect + client.on_message = on_message + client.reconnect_delay_set(min_delay=1, max_delay=30) + try: + client.connect(MQTT_BROKER_URL, MQTT_BROKER_PORT) + except Exception as e: + connection_state = "reconnecting" + print(f"[MQTT] Initial connect failed: {e} — will keep retrying in background") + client.loop_start() + return client + +def get_connection_state(): + return connection_state + +def get_latest_events(): + return list(latest_events.values()) \ No newline at end of file diff --git a/src/production/backend/backend/project-echo-openapi.json b/src/production/backend/backend/project-echo-openapi.json index 6c863b906..7c41baf19 100644 --- a/src/production/backend/backend/project-echo-openapi.json +++ b/src/production/backend/backend/project-echo-openapi.json @@ -6,6 +6,25 @@ "version": "1.0.0" }, "paths": { + "/mqtt/connection-state": { + "get": { + "tags": [ + "mqtt" + ], + "summary": "Mqtt Connection State", + "operationId": "mqtt_connection_state_mqtt_connection_state_get", + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": {} + } + } + } + } + } + }, "/api/audio/upload": { "post": { "tags": [ diff --git a/src/production/hmi/ui/public/js/HMI.js b/src/production/hmi/ui/public/js/HMI.js index a19c371a8..03bba50bf 100644 --- a/src/production/hmi/ui/public/js/HMI.js +++ b/src/production/hmi/ui/public/js/HMI.js @@ -12,7 +12,7 @@ * No changes to map logic, layer management, or audio handling. */ -import { showToast, getApiErrorMessage, withRetry } from "./HMI-utils.js"; +import { showToast, getApiErrorMessage, withRetry, showPageBanner, hidePageBanner } from "./HMI-utils.js"; import { getAudioRecorder } from "./audio_recorder.js"; import { AudioDecoder, @@ -590,8 +590,118 @@ fetch("./js/sample_data.json") * Task 7: error message now routed through getApiErrorMessage so the wording * is consistent with every other error surface in the application. */ +// ───────────────────────────────────────────────────────────────────────────── +// MQTT connection state polling (FR-A2) +// ───────────────────────────────────────────────────────────────────────────── + +let _lastMqttState = null; + +async function pollMqttConnectionState() { + try { + const response = await fetch("http://localhost:9000/mqtt/connection-state"); + if (!response.ok) throw new Error("Failed to fetch connection state"); + const data = await response.json(); + const state = data.state; + + if (state !== _lastMqttState) { + if (state === "connected") { + hidePageBanner("warning"); + hidePageBanner("error"); + if (_lastMqttState !== null) { + showToast("Live data connection restored", "success"); + } + } else if (state === "reconnecting") { + showPageBanner("Live data connection lost — reconnecting…", "warning", false); + showToast("Live data connection lost, reconnecting…", "warning"); + } else if (state === "disconnected") { + showPageBanner("Live data unavailable", "error", false); + } + _lastMqttState = state; + } + } catch (err) { + console.error("Error polling MQTT connection state:", err); + + // The status check itself failed (backend unreachable) — treat this + // as unavailable rather than silently keeping the last-known state. + if (_lastMqttState !== "unavailable") { + showPageBanner("Live data unavailable — unable to check connection status", "error", false); + _lastMqttState = "unavailable"; + } + } +} + +function startMqttConnectionPolling() { + pollMqttConnectionState(); + setInterval(pollMqttConnectionState, 5000); +} + +// ───────────────────────────────────────────────────────────────────────────── +// Live MQTT events → map layers (FR-A2) +// ───────────────────────────────────────────────────────────────────────────── + +const _seenMqttEventIds = new Set(); + +async function pollMqttLatestEvents(hmiState) { + try { + const response = await fetch("http://localhost:9000/mqtt/latest-events"); + if (!response.ok) throw new Error("Failed to fetch latest events"); + const data = await response.json(); + const events = data.events || []; + + + const newVocalizationEvents = []; + const newMovementEvents = []; + + for (const event of events) { + if (_seenMqttEventIds.has(event._id)) continue; + _seenMqttEventIds.add(event._id); + + switch (event.eventType) { + case "vocalization": + newVocalizationEvents.push(event); + break; + case "movement": + newMovementEvents.push(event); + break; + case "sensor_health": + case "iot_node": + document.dispatchEvent( + new CustomEvent(`mqtt:${event.eventType}`, { detail: event }) + ); + break; + } + + + } + + if (newVocalizationEvents.length > 0) { + updateVocalizationLayerFromLiveData(hmiState, newVocalizationEvents); + } + if (newMovementEvents.length > 0) { + updateAnimalMovementLayerFromLiveData(hmiState, newMovementEvents); + } + + + } catch (err) { + console.error("Error polling MQTT latest events:", err); + } +} + +function startMqttEventPolling(hmiState) { + pollMqttLatestEvents(hmiState); + setInterval(() => pollMqttLatestEvents(hmiState), 5000); +} + + + + + + + export function initialiseHMI(hmiState) { console.log("initialising"); + startMqttConnectionPolling(); + startMqttEventPolling(hmiState); showMapSpinner("Loading map data…"); hideMapError(); diff --git a/src/production/simulator/src/comms_manager.py b/src/production/simulator/src/comms_manager.py index 4a13d6af0..db33b5ad0 100644 --- a/src/production/simulator/src/comms_manager.py +++ b/src/production/simulator/src/comms_manager.py @@ -162,6 +162,11 @@ def echo_api_send_animal_movement(self, animal): "animalTrueLLA": list(animal.getLLA()) } + self.mqtt_client.publish( + os.environ.get("MQTT_MOVEMENT_TOPIC", "projectecho/movement"), + json.dumps(movement_event), + ) + url = 'http://ts-api-cont:9000/sim/movement' x = requests.post(url, json = movement_event)