-
Notifications
You must be signed in to change notification settings - Fork 34
EH/KIM/fr-a2-mqtt-event-handling #961
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
a2f3bd8
14cf37c
664685b
4965635
215ca17
2f29aeb
b1e6261
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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 | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This dict just keeps growing, every message that comes in adds a new key and nothing ever removes one. Over a long-running deployment that's a slow memory leak on the backend, and /mqtt/latest-events will keep returning a bigger response every time since it's returning everything ever seen, not just recent stuff. Also, this write happens on the MQTT client's own thread (from loop_start()), while the /mqtt/latest-events handler reads it from a different thread with no lock in between. list(latest_events.values()) while this line is mutating the dict on another thread could occasionally throw a "dictionary changed size during iteration" error. Would suggest capping this (e.g. a bounded deque or evicting old entries) and adding a lock around the read/write, similar to what's already being done for connection_state elsewhere maybe not, but worth double checking that one too. |
||
| # TODO: forward `normalized` to wherever the frontend/dashboard reads from | ||
|
kimmynoto marked this conversation as resolved.
|
||
|
|
||
| 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: | ||
|
AnDo081105 marked this conversation as resolved.
|
||
| 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()) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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"; | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This whole commit (Start FR-A4: Live Map Controls and Filters) landed after Khanh's approval, so it hasn't actually been reviewed, and it's out of scope for this PR, this one's about FR-A2/MQTT. Can we split the FR-A4 map controls work into its own PR, and either fix the syntax error in audio_recorder.js first or leave this import commented out until it's actually fixed? Don't think this should go in as part of the MQTT work. |
||
| 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"); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This is hardcoded to http://localhost:9000, same thing a couple lines down in pollMqttLatestEvents for /mqtt/latest-events. Works fine for the local docker-compose setup but will just fail for anyone running this anywhere else, since there's no way to point it at a different host. Might be worth pulling this from a config value or the same base URL the rest of the app's API calls already use, rather than hardcoding it twice here. |
||
| 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); | ||
|
kimmynoto marked this conversation as resolved.
|
||
|
|
||
| // 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(); | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.