Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion src/deployment/docker/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
28 changes: 28 additions & 0 deletions src/production/backend/app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Comment thread
AnDo081105 marked this conversation as resolved.
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')
Expand Down
28 changes: 28 additions & 0 deletions src/production/backend/app/routers/hmi.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 = [
Expand Down
133 changes: 133 additions & 0 deletions src/production/backend/app/services/mqtt_client.py
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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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
Comment thread
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:
Comment thread
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())
19 changes: 19 additions & 0 deletions src/production/backend/backend/project-echo-openapi.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": [
Expand Down
112 changes: 111 additions & 1 deletion src/production/hmi/ui/public/js/HMI.js
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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.
Bigger issue though: you're re-enabling this import with a comment saying audio_recorder.js has a pre-existing syntax error. If that's accurate, importing it will throw and could break the whole HMI page on load, not just the new map controls.

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,
Expand Down Expand Up @@ -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");

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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);
Comment thread
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();
Expand Down
5 changes: 5 additions & 0 deletions src/production/simulator/src/comms_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading