diff --git a/CHANGELOG.md b/CHANGELOG.md
index 338c0697a..9e98bb3bc 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -8,8 +8,10 @@ All notable changes to Bambuddy will be documented in this file.
- **MQTT Smart Plug Support** - Add smart plugs that subscribe to MQTT topics for energy monitoring (Issue #173):
- New "MQTT" plug type alongside Tasmota and Home Assistant
- Subscribe to any MQTT topic (Zigbee2MQTT, Shelly, Tasmota discovery, etc.)
- - Configurable JSON paths for power, energy, and state extraction (e.g., `power_l1`, `data.power`)
- - Optional multiplier for unit conversion (mW to W, etc.)
+ - **Separate topics per data type**: Configure different MQTT topics for power, energy, and state
+ - Configurable JSON paths for data extraction (e.g., `power_l1`, `data.power`)
+ - **Separate multipliers**: Individual multiplier for power and energy (e.g., mW→W, Wh→kWh)
+ - **Custom ON value**: Configure what value means "ON" for state (e.g., "ON", "true", "1")
- Monitor-only: displays power/energy data without control capabilities
- Reuses existing MQTT broker settings from Settings → Network
- Energy data included in statistics and per-print tracking
diff --git a/README.md b/README.md
index 999972dbb..6d42769e1 100644
--- a/README.md
+++ b/README.md
@@ -78,7 +78,8 @@
- Per-printer AMS mapping (individual slot configuration for print farms)
- Scheduled prints (date/time)
- Queue Only mode (stage without auto-start)
-- Smart plug integration (Tasmota, Home Assistant)
+- Smart plug integration (Tasmota, Home Assistant, MQTT)
+- MQTT smart plugs: Subscribe to Zigbee2MQTT, Shelly, or any MQTT topic for energy monitoring
- Energy consumption tracking (per-print kWh and cost)
- HA energy sensor support (for plugs with separate power/energy sensors)
- Auto power-on before print
diff --git a/backend/app/api/routes/settings.py b/backend/app/api/routes/settings.py
index 798b343b8..67f2fd675 100644
--- a/backend/app/api/routes/settings.py
+++ b/backend/app/api/routes/settings.py
@@ -360,12 +360,21 @@ async def export_backup(
"ha_power_entity": plug.ha_power_entity,
"ha_energy_today_entity": plug.ha_energy_today_entity,
"ha_energy_total_entity": plug.ha_energy_total_entity,
- # MQTT plug fields
+ # MQTT plug fields (legacy)
"mqtt_topic": plug.mqtt_topic,
- "mqtt_power_path": plug.mqtt_power_path,
- "mqtt_energy_path": plug.mqtt_energy_path,
- "mqtt_state_path": plug.mqtt_state_path,
"mqtt_multiplier": plug.mqtt_multiplier,
+ # MQTT power fields
+ "mqtt_power_topic": plug.mqtt_power_topic,
+ "mqtt_power_path": plug.mqtt_power_path,
+ "mqtt_power_multiplier": plug.mqtt_power_multiplier,
+ # MQTT energy fields
+ "mqtt_energy_topic": plug.mqtt_energy_topic,
+ "mqtt_energy_path": plug.mqtt_energy_path,
+ "mqtt_energy_multiplier": plug.mqtt_energy_multiplier,
+ # MQTT state fields
+ "mqtt_state_topic": plug.mqtt_state_topic,
+ "mqtt_state_path": plug.mqtt_state_path,
+ "mqtt_state_on_value": plug.mqtt_state_on_value,
"printer_serial": printer_id_to_serial.get(plug.printer_id) if plug.printer_id else None,
"enabled": plug.enabled,
"auto_on": plug.auto_on,
@@ -1270,10 +1279,16 @@ async def import_backup(
result = await db.execute(select(SmartPlug).where(SmartPlug.ha_entity_id == plug_data["ha_entity_id"]))
existing = result.scalar_one_or_none()
plug_identifier = plug_data["ha_entity_id"]
- elif plug_type == "mqtt" and plug_data.get("mqtt_topic"):
- result = await db.execute(select(SmartPlug).where(SmartPlug.mqtt_topic == plug_data["mqtt_topic"]))
+ elif plug_type == "mqtt" and (plug_data.get("mqtt_power_topic") or plug_data.get("mqtt_topic")):
+ # Check by mqtt_power_topic first (new format), fall back to mqtt_topic (legacy)
+ power_topic = plug_data.get("mqtt_power_topic") or plug_data.get("mqtt_topic")
+ result = await db.execute(
+ select(SmartPlug).where(
+ (SmartPlug.mqtt_power_topic == power_topic) | (SmartPlug.mqtt_topic == power_topic)
+ )
+ )
existing = result.scalar_one_or_none()
- plug_identifier = plug_data["mqtt_topic"]
+ plug_identifier = power_topic
elif plug_data.get("ip_address"):
result = await db.execute(select(SmartPlug).where(SmartPlug.ip_address == plug_data["ip_address"]))
existing = result.scalar_one_or_none()
@@ -1290,12 +1305,21 @@ async def import_backup(
existing.ha_power_entity = plug_data.get("ha_power_entity")
existing.ha_energy_today_entity = plug_data.get("ha_energy_today_entity")
existing.ha_energy_total_entity = plug_data.get("ha_energy_total_entity")
- # MQTT fields
+ # MQTT fields (legacy)
existing.mqtt_topic = plug_data.get("mqtt_topic")
- existing.mqtt_power_path = plug_data.get("mqtt_power_path")
- existing.mqtt_energy_path = plug_data.get("mqtt_energy_path")
- existing.mqtt_state_path = plug_data.get("mqtt_state_path")
existing.mqtt_multiplier = plug_data.get("mqtt_multiplier", 1.0)
+ # MQTT power fields
+ existing.mqtt_power_topic = plug_data.get("mqtt_power_topic")
+ existing.mqtt_power_path = plug_data.get("mqtt_power_path")
+ existing.mqtt_power_multiplier = plug_data.get("mqtt_power_multiplier", 1.0)
+ # MQTT energy fields
+ existing.mqtt_energy_topic = plug_data.get("mqtt_energy_topic")
+ existing.mqtt_energy_path = plug_data.get("mqtt_energy_path")
+ existing.mqtt_energy_multiplier = plug_data.get("mqtt_energy_multiplier", 1.0)
+ # MQTT state fields
+ existing.mqtt_state_topic = plug_data.get("mqtt_state_topic")
+ existing.mqtt_state_path = plug_data.get("mqtt_state_path")
+ existing.mqtt_state_on_value = plug_data.get("mqtt_state_on_value")
existing.printer_id = printer_id
existing.enabled = plug_data.get("enabled", True)
existing.auto_on = plug_data.get("auto_on", True)
@@ -1325,12 +1349,21 @@ async def import_backup(
ha_power_entity=plug_data.get("ha_power_entity"),
ha_energy_today_entity=plug_data.get("ha_energy_today_entity"),
ha_energy_total_entity=plug_data.get("ha_energy_total_entity"),
- # MQTT fields
+ # MQTT fields (legacy)
mqtt_topic=plug_data.get("mqtt_topic"),
- mqtt_power_path=plug_data.get("mqtt_power_path"),
- mqtt_energy_path=plug_data.get("mqtt_energy_path"),
- mqtt_state_path=plug_data.get("mqtt_state_path"),
mqtt_multiplier=plug_data.get("mqtt_multiplier", 1.0),
+ # MQTT power fields
+ mqtt_power_topic=plug_data.get("mqtt_power_topic"),
+ mqtt_power_path=plug_data.get("mqtt_power_path"),
+ mqtt_power_multiplier=plug_data.get("mqtt_power_multiplier", 1.0),
+ # MQTT energy fields
+ mqtt_energy_topic=plug_data.get("mqtt_energy_topic"),
+ mqtt_energy_path=plug_data.get("mqtt_energy_path"),
+ mqtt_energy_multiplier=plug_data.get("mqtt_energy_multiplier", 1.0),
+ # MQTT state fields
+ mqtt_state_topic=plug_data.get("mqtt_state_topic"),
+ mqtt_state_path=plug_data.get("mqtt_state_path"),
+ mqtt_state_on_value=plug_data.get("mqtt_state_on_value"),
printer_id=printer_id,
enabled=plug_data.get("enabled", True),
auto_on=plug_data.get("auto_on", True),
diff --git a/backend/app/api/routes/smart_plugs.py b/backend/app/api/routes/smart_plugs.py
index b782e126c..42c915347 100644
--- a/backend/app/api/routes/smart_plugs.py
+++ b/backend/app/api/routes/smart_plugs.py
@@ -96,17 +96,44 @@ async def create_smart_plug(
await db.commit()
await db.refresh(plug)
- # Subscribe MQTT plugs to their topic
- if plug.plug_type == "mqtt" and plug.mqtt_topic:
- mqtt_relay.smart_plug_service.subscribe(
- plug_id=plug.id,
- topic=plug.mqtt_topic,
- power_path=plug.mqtt_power_path,
- energy_path=plug.mqtt_energy_path,
- state_path=plug.mqtt_state_path,
- multiplier=plug.mqtt_multiplier or 1.0,
- )
- logger.info(f"Created MQTT plug '{plug.name}' subscribed to {plug.mqtt_topic}")
+ # Subscribe MQTT plugs to their topics
+ if plug.plug_type == "mqtt":
+ # Determine effective topics (new fields take priority, fall back to legacy)
+ power_topic = plug.mqtt_power_topic or plug.mqtt_topic
+ energy_topic = plug.mqtt_energy_topic or plug.mqtt_topic
+ state_topic = plug.mqtt_state_topic or plug.mqtt_topic
+
+ # Only subscribe if at least one data source is configured
+ if (
+ (power_topic and plug.mqtt_power_path)
+ or (energy_topic and plug.mqtt_energy_path)
+ or (state_topic and plug.mqtt_state_path)
+ ):
+ mqtt_relay.smart_plug_service.subscribe(
+ plug_id=plug.id,
+ # Power source
+ power_topic=power_topic if plug.mqtt_power_path else None,
+ power_path=plug.mqtt_power_path,
+ power_multiplier=plug.mqtt_power_multiplier or plug.mqtt_multiplier or 1.0,
+ # Energy source
+ energy_topic=energy_topic if plug.mqtt_energy_path else None,
+ energy_path=plug.mqtt_energy_path,
+ energy_multiplier=plug.mqtt_energy_multiplier or plug.mqtt_multiplier or 1.0,
+ # State source
+ state_topic=state_topic if plug.mqtt_state_path else None,
+ state_path=plug.mqtt_state_path,
+ state_on_value=plug.mqtt_state_on_value,
+ )
+ topics = [
+ t
+ for t in [
+ power_topic if plug.mqtt_power_path else None,
+ energy_topic if plug.mqtt_energy_path else None,
+ state_topic if plug.mqtt_state_path else None,
+ ]
+ if t
+ ]
+ logger.info(f"Created MQTT plug '{plug.name}' subscribed to {', '.join(set(topics))}")
elif plug.plug_type == "homeassistant":
logger.info(f"Created Home Assistant plug '{plug.name}' ({plug.ha_entity_id})")
else:
@@ -338,9 +365,19 @@ async def update_smart_plug(
if result.scalar_one_or_none():
raise HTTPException(400, "This printer already has a smart plug assigned")
- # Check if MQTT topic is changing - need to resubscribe
- old_topic = plug.mqtt_topic
+ # Track old MQTT settings for comparison
old_plug_type = plug.plug_type
+ old_mqtt_config = {
+ "power_topic": plug.mqtt_power_topic or plug.mqtt_topic,
+ "power_path": plug.mqtt_power_path,
+ "power_multiplier": plug.mqtt_power_multiplier,
+ "energy_topic": plug.mqtt_energy_topic or plug.mqtt_topic,
+ "energy_path": plug.mqtt_energy_path,
+ "energy_multiplier": plug.mqtt_energy_multiplier,
+ "state_topic": plug.mqtt_state_topic or plug.mqtt_topic,
+ "state_path": plug.mqtt_state_path,
+ "state_on_value": plug.mqtt_state_on_value,
+ }
for field, value in update_data.items():
setattr(plug, field, value)
@@ -353,20 +390,50 @@ async def update_smart_plug(
# Changed away from MQTT - unsubscribe
mqtt_relay.smart_plug_service.unsubscribe(plug.id)
elif plug.plug_type == "mqtt":
- # Is now MQTT - check if topic changed or newly MQTT
- if old_plug_type != "mqtt" or old_topic != plug.mqtt_topic:
- # Unsubscribe from old topic first
+ # Check if any MQTT config changed
+ new_mqtt_config = {
+ "power_topic": plug.mqtt_power_topic or plug.mqtt_topic,
+ "power_path": plug.mqtt_power_path,
+ "power_multiplier": plug.mqtt_power_multiplier,
+ "energy_topic": plug.mqtt_energy_topic or plug.mqtt_topic,
+ "energy_path": plug.mqtt_energy_path,
+ "energy_multiplier": plug.mqtt_energy_multiplier,
+ "state_topic": plug.mqtt_state_topic or plug.mqtt_topic,
+ "state_path": plug.mqtt_state_path,
+ "state_on_value": plug.mqtt_state_on_value,
+ }
+
+ mqtt_changed = old_plug_type != "mqtt" or old_mqtt_config != new_mqtt_config
+
+ if mqtt_changed:
+ # Unsubscribe from old topics first
if old_plug_type == "mqtt":
mqtt_relay.smart_plug_service.unsubscribe(plug.id)
- # Subscribe to new topic
- if plug.mqtt_topic:
+
+ # Subscribe to new topics
+ power_topic = plug.mqtt_power_topic or plug.mqtt_topic
+ energy_topic = plug.mqtt_energy_topic or plug.mqtt_topic
+ state_topic = plug.mqtt_state_topic or plug.mqtt_topic
+
+ if (
+ (power_topic and plug.mqtt_power_path)
+ or (energy_topic and plug.mqtt_energy_path)
+ or (state_topic and plug.mqtt_state_path)
+ ):
mqtt_relay.smart_plug_service.subscribe(
plug_id=plug.id,
- topic=plug.mqtt_topic,
+ # Power source
+ power_topic=power_topic if plug.mqtt_power_path else None,
power_path=plug.mqtt_power_path,
+ power_multiplier=plug.mqtt_power_multiplier or plug.mqtt_multiplier or 1.0,
+ # Energy source
+ energy_topic=energy_topic if plug.mqtt_energy_path else None,
energy_path=plug.mqtt_energy_path,
+ energy_multiplier=plug.mqtt_energy_multiplier or plug.mqtt_multiplier or 1.0,
+ # State source
+ state_topic=state_topic if plug.mqtt_state_path else None,
state_path=plug.mqtt_state_path,
- multiplier=plug.mqtt_multiplier or 1.0,
+ state_on_value=plug.mqtt_state_on_value,
)
logger.info(f"Updated smart plug '{plug.name}'")
diff --git a/backend/app/core/database.py b/backend/app/core/database.py
index 683347868..9b3c6c299 100644
--- a/backend/app/core/database.py
+++ b/backend/app/core/database.py
@@ -760,7 +760,7 @@ async def run_migrations(conn):
except Exception:
pass
- # Migration: Add MQTT smart plug fields
+ # Migration: Add MQTT smart plug fields (legacy)
try:
await conn.execute(text("ALTER TABLE smart_plugs ADD COLUMN mqtt_topic VARCHAR(200)"))
except Exception:
@@ -782,6 +782,45 @@ async def run_migrations(conn):
except Exception:
pass
+ # Migration: Add enhanced MQTT smart plug fields (separate topics and multipliers)
+ try:
+ await conn.execute(text("ALTER TABLE smart_plugs ADD COLUMN mqtt_power_topic VARCHAR(200)"))
+ except Exception:
+ pass
+ try:
+ await conn.execute(text("ALTER TABLE smart_plugs ADD COLUMN mqtt_power_multiplier REAL DEFAULT 1.0"))
+ except Exception:
+ pass
+ try:
+ await conn.execute(text("ALTER TABLE smart_plugs ADD COLUMN mqtt_energy_topic VARCHAR(200)"))
+ except Exception:
+ pass
+ try:
+ await conn.execute(text("ALTER TABLE smart_plugs ADD COLUMN mqtt_energy_multiplier REAL DEFAULT 1.0"))
+ except Exception:
+ pass
+ try:
+ await conn.execute(text("ALTER TABLE smart_plugs ADD COLUMN mqtt_state_topic VARCHAR(200)"))
+ except Exception:
+ pass
+ try:
+ await conn.execute(text("ALTER TABLE smart_plugs ADD COLUMN mqtt_state_on_value VARCHAR(50)"))
+ except Exception:
+ pass
+
+ # Migration: Copy existing mqtt_topic to mqtt_power_topic for backward compatibility
+ try:
+ await conn.execute(
+ text("""
+ UPDATE smart_plugs
+ SET mqtt_power_topic = mqtt_topic,
+ mqtt_power_multiplier = mqtt_multiplier
+ WHERE mqtt_topic IS NOT NULL AND mqtt_power_topic IS NULL
+ """)
+ )
+ except Exception:
+ pass
+
async def seed_notification_templates():
"""Seed default notification templates if they don't exist."""
diff --git a/backend/app/models/smart_plug.py b/backend/app/models/smart_plug.py
index 99d63ed87..3bcc3b1ee 100644
--- a/backend/app/models/smart_plug.py
+++ b/backend/app/models/smart_plug.py
@@ -25,13 +25,30 @@ class SmartPlug(Base):
ha_energy_total_entity: Mapped[str | None] = mapped_column(String(100), nullable=True) # sensor.xxx_total
# MQTT plug fields (required when plug_type="mqtt")
+ # Legacy field - kept for backward compatibility, now use mqtt_power_topic
mqtt_topic: Mapped[str | None] = mapped_column(
String(200), nullable=True
- ) # e.g., "zigbee2mqtt/shelly-working-room"
+ ) # e.g., "zigbee2mqtt/shelly-working-room" (deprecated, use mqtt_power_topic)
+
+ # Power monitoring
+ mqtt_power_topic: Mapped[str | None] = mapped_column(String(200), nullable=True) # Topic for power data
mqtt_power_path: Mapped[str | None] = mapped_column(String(100), nullable=True) # e.g., "power_l1" or "data.power"
+ mqtt_power_multiplier: Mapped[float] = mapped_column(Float, default=1.0) # Unit conversion for power
+
+ # Energy monitoring
+ mqtt_energy_topic: Mapped[str | None] = mapped_column(String(200), nullable=True) # Topic for energy data
mqtt_energy_path: Mapped[str | None] = mapped_column(String(100), nullable=True) # e.g., "energy_l1"
+ mqtt_energy_multiplier: Mapped[float] = mapped_column(Float, default=1.0) # Unit conversion for energy
+
+ # State monitoring
+ mqtt_state_topic: Mapped[str | None] = mapped_column(String(200), nullable=True) # Topic for state data
mqtt_state_path: Mapped[str | None] = mapped_column(String(100), nullable=True) # e.g., "state_l1" for ON/OFF
- mqtt_multiplier: Mapped[float] = mapped_column(Float, default=1.0) # Unit conversion (e.g., 0.001 for mW→W)
+ mqtt_state_on_value: Mapped[str | None] = mapped_column(
+ String(50), nullable=True
+ ) # What value means "ON" (e.g., "ON", "true", "1")
+
+ # Legacy multiplier - kept for backward compatibility
+ mqtt_multiplier: Mapped[float] = mapped_column(Float, default=1.0) # Deprecated, use mqtt_power_multiplier
# Link to printer (1:1)
printer_id: Mapped[int | None] = mapped_column(
diff --git a/backend/app/schemas/smart_plug.py b/backend/app/schemas/smart_plug.py
index ed4d2be5c..edb8a1210 100644
--- a/backend/app/schemas/smart_plug.py
+++ b/backend/app/schemas/smart_plug.py
@@ -21,11 +21,28 @@ class SmartPlugBase(BaseModel):
ha_energy_total_entity: str | None = Field(default=None, pattern=r"^sensor\.[a-z0-9_]+$")
# MQTT fields (required when plug_type="mqtt")
- mqtt_topic: str | None = Field(default=None, max_length=200)
+ # Legacy field - kept for backward compatibility
+ mqtt_topic: str | None = Field(default=None, max_length=200) # Deprecated, use mqtt_power_topic
+
+ # Power monitoring
+ mqtt_power_topic: str | None = Field(default=None, max_length=200) # Topic for power data
mqtt_power_path: str | None = Field(default=None, max_length=100) # e.g., "power_l1" or "data.power"
+ mqtt_power_multiplier: float = Field(default=1.0, ge=0.0001, le=10000) # Unit conversion for power
+
+ # Energy monitoring
+ mqtt_energy_topic: str | None = Field(default=None, max_length=200) # Topic for energy data
mqtt_energy_path: str | None = Field(default=None, max_length=100) # e.g., "energy_l1"
+ mqtt_energy_multiplier: float = Field(default=1.0, ge=0.0001, le=10000) # Unit conversion for energy
+
+ # State monitoring
+ mqtt_state_topic: str | None = Field(default=None, max_length=200) # Topic for state data
mqtt_state_path: str | None = Field(default=None, max_length=100) # e.g., "state_l1" for ON/OFF
- mqtt_multiplier: float = Field(default=1.0, ge=0.0001, le=10000) # Unit conversion (e.g., 0.001 for mW→W)
+ mqtt_state_on_value: str | None = Field(
+ default=None, max_length=50
+ ) # What value means "ON" (e.g., "ON", "true", "1")
+
+ # Legacy multiplier - kept for backward compatibility
+ mqtt_multiplier: float = Field(default=1.0, ge=0.0001, le=10000) # Deprecated, use mqtt_power_multiplier
printer_id: int | None = None
enabled: bool = True
@@ -52,10 +69,18 @@ class SmartPlugBase(BaseModel):
if self.plug_type == "homeassistant" and not self.ha_entity_id:
raise ValueError("ha_entity_id is required for Home Assistant plugs")
if self.plug_type == "mqtt":
- if not self.mqtt_topic:
- raise ValueError("mqtt_topic is required for MQTT plugs")
- if not self.mqtt_power_path and not self.mqtt_state_path:
- raise ValueError("At least mqtt_power_path or mqtt_state_path is required for MQTT plugs")
+ # Determine the effective power topic (new field takes priority, fall back to legacy)
+ power_topic = self.mqtt_power_topic or self.mqtt_topic
+ has_power = power_topic and self.mqtt_power_path
+ has_energy = self.mqtt_energy_topic and self.mqtt_energy_path
+ has_state = self.mqtt_state_topic and self.mqtt_state_path
+
+ # At least one data source must be fully configured
+ if not has_power and not has_energy and not has_state:
+ raise ValueError(
+ "At least one MQTT data source must be configured: "
+ "power (topic + path), energy (topic + path), or state (topic + path)"
+ )
return self
@@ -72,12 +97,21 @@ class SmartPlugUpdate(BaseModel):
ha_power_entity: str | None = None
ha_energy_today_entity: str | None = None
ha_energy_total_entity: str | None = None
- # MQTT fields
+ # MQTT fields (legacy)
mqtt_topic: str | None = None
- mqtt_power_path: str | None = None
- mqtt_energy_path: str | None = None
- mqtt_state_path: str | None = None
mqtt_multiplier: float | None = Field(default=None, ge=0.0001, le=10000)
+ # MQTT power fields
+ mqtt_power_topic: str | None = None
+ mqtt_power_path: str | None = None
+ mqtt_power_multiplier: float | None = Field(default=None, ge=0.0001, le=10000)
+ # MQTT energy fields
+ mqtt_energy_topic: str | None = None
+ mqtt_energy_path: str | None = None
+ mqtt_energy_multiplier: float | None = Field(default=None, ge=0.0001, le=10000)
+ # MQTT state fields
+ mqtt_state_topic: str | None = None
+ mqtt_state_path: str | None = None
+ mqtt_state_on_value: str | None = None
printer_id: int | None = None
enabled: bool | None = None
auto_on: bool | None = None
diff --git a/backend/app/services/mqtt_smart_plug.py b/backend/app/services/mqtt_smart_plug.py
index 35e913153..dfc7fbba0 100644
--- a/backend/app/services/mqtt_smart_plug.py
+++ b/backend/app/services/mqtt_smart_plug.py
@@ -26,6 +26,16 @@ class SmartPlugMQTTData:
last_seen: datetime = field(default_factory=datetime.utcnow)
+@dataclass
+class MQTTDataSourceConfig:
+ """Configuration for a single MQTT data source (power, energy, or state)."""
+
+ topic: str
+ path: str
+ multiplier: float = 1.0 # For power/energy
+ on_value: str | None = None # For state (what value means "ON")
+
+
class MQTTSmartPlugService:
"""Subscribes to MQTT topics for smart plug energy monitoring."""
@@ -36,10 +46,10 @@ class MQTTSmartPlugService:
self.client: mqtt.Client | None = None
self.connected = False
self._lock = threading.Lock()
- # topic -> list of plug_ids (multiple plugs can subscribe to same topic with different paths)
- self.subscriptions: dict[str, list[int]] = {}
- # plug_id -> (topic, power_path, energy_path, state_path, multiplier)
- self.plug_configs: dict[int, tuple[str, str | None, str | None, str | None, float]] = {}
+ # topic -> list of (plug_id, data_type) where data_type is "power", "energy", or "state"
+ self.subscriptions: dict[str, list[tuple[int, str]]] = {}
+ # plug_id -> {data_type: MQTTDataSourceConfig}
+ self.plug_configs: dict[int, dict[str, MQTTDataSourceConfig]] = {}
# plug_id -> latest data
self.plug_data: dict[int, SmartPlugMQTTData] = {}
self._configured = False
@@ -205,78 +215,80 @@ class MQTTSmartPlugService:
topic = msg.topic
with self._lock:
- plug_ids = self.subscriptions.get(topic, [])
- if not plug_ids:
+ subscriptions = self.subscriptions.get(topic, [])
+ if not subscriptions:
return
- # Parse JSON payload
+ # Parse JSON payload (or treat as raw value)
try:
payload = json.loads(msg.payload.decode("utf-8"))
- except (json.JSONDecodeError, UnicodeDecodeError) as e:
- logger.debug(f"MQTT smart plug: failed to parse message on {topic}: {e}")
- return
+ is_json = True
+ except (json.JSONDecodeError, UnicodeDecodeError):
+ # Not JSON - treat the whole payload as a raw value
+ payload = msg.payload.decode("utf-8").strip()
+ is_json = False
- # Process for each subscribed plug
- for plug_id in plug_ids:
- config = self.plug_configs.get(plug_id)
+ # Process for each subscribed (plug_id, data_type)
+ for plug_id, data_type in subscriptions:
+ configs = self.plug_configs.get(plug_id, {})
+ config = configs.get(data_type)
if not config:
continue
- _, power_path, energy_path, state_path, multiplier = config
-
- # Extract values
- power = None
- energy = None
- state = None
-
- if power_path:
- raw_power = self._extract_json_path(payload, power_path)
- if raw_power is not None:
- try:
- power = float(raw_power) * multiplier
- except (ValueError, TypeError):
- pass
-
- if energy_path:
- raw_energy = self._extract_json_path(payload, energy_path)
- if raw_energy is not None:
- try:
- energy = float(raw_energy) * multiplier
- except (ValueError, TypeError):
- pass
-
- if state_path:
- raw_state = self._extract_json_path(payload, state_path)
- if raw_state is not None:
- # Normalize state to ON/OFF
- state_str = str(raw_state).upper()
- if state_str in ("ON", "1", "TRUE"):
- state = "ON"
- elif state_str in ("OFF", "0", "FALSE"):
- state = "OFF"
- else:
- state = state_str
-
- # Update plug data
- if plug_id in self.plug_data:
- data = self.plug_data[plug_id]
- if power is not None:
- data.power = power
- if energy is not None:
- data.energy = energy
- if state is not None:
- data.state = state
- data.last_seen = datetime.utcnow()
+ # Extract value using path (or use raw payload if no path or not JSON)
+ if is_json and config.path:
+ raw_value = self._extract_json_path(payload, config.path)
+ elif is_json and not config.path:
+ # JSON but no path - use the whole payload (shouldn't happen normally)
+ raw_value = payload
else:
- self.plug_data[plug_id] = SmartPlugMQTTData(
- plug_id=plug_id,
- power=power,
- energy=energy,
- state=state,
- last_seen=datetime.utcnow(),
- )
+ # Raw value (non-JSON)
+ raw_value = payload
- logger.debug(f"MQTT smart plug {plug_id}: power={power}, energy={energy}, state={state}")
+ if raw_value is None:
+ continue
+
+ # Initialize plug data if needed
+ if plug_id not in self.plug_data:
+ self.plug_data[plug_id] = SmartPlugMQTTData(plug_id=plug_id)
+
+ data = self.plug_data[plug_id]
+ data.last_seen = datetime.utcnow()
+
+ # Process based on data type
+ if data_type == "power":
+ try:
+ data.power = float(raw_value) * config.multiplier
+ logger.debug(f"MQTT smart plug {plug_id}: power={data.power}")
+ except (ValueError, TypeError):
+ pass
+
+ elif data_type == "energy":
+ try:
+ data.energy = float(raw_value) * config.multiplier
+ logger.debug(f"MQTT smart plug {plug_id}: energy={data.energy}")
+ except (ValueError, TypeError):
+ pass
+
+ elif data_type == "state":
+ state_str = str(raw_value)
+ # Check against configured ON value if set
+ if config.on_value:
+ # Case-insensitive comparison
+ if state_str.lower() == config.on_value.lower():
+ data.state = "ON"
+ else:
+ data.state = "OFF"
+ else:
+ # Default behavior: normalize common values
+ upper_state = state_str.upper()
+ if upper_state in ("ON", "1", "TRUE"):
+ data.state = "ON"
+ elif upper_state in ("OFF", "0", "FALSE"):
+ data.state = "OFF"
+ else:
+ data.state = state_str
+ logger.debug(f"MQTT smart plug {plug_id}: state={data.state}")
def _extract_json_path(self, data: dict, path: str) -> Any:
"""Extract value using dot notation (e.g., 'power_l1' or 'data.power').
@@ -304,61 +316,131 @@ class MQTTSmartPlugService:
with self._lock:
for topic in self.subscriptions:
- try:
- self.client.subscribe(topic, qos=1)
- logger.debug(f"MQTT smart plug: resubscribed to {topic}")
- except Exception as e:
- logger.error(f"MQTT smart plug: failed to resubscribe to {topic}: {e}")
+ if self.subscriptions[topic]: # Only if there are subscribers
+ try:
+ self.client.subscribe(topic, qos=1)
+ logger.debug(f"MQTT smart plug: resubscribed to {topic}")
+ except Exception as e:
+ logger.error(f"MQTT smart plug: failed to resubscribe to {topic}: {e}")
def subscribe(
self,
plug_id: int,
- topic: str,
+ # Power source
+ power_topic: str | None = None,
power_path: str | None = None,
+ power_multiplier: float = 1.0,
+ # Energy source
+ energy_topic: str | None = None,
energy_path: str | None = None,
+ energy_multiplier: float = 1.0,
+ # State source
+ state_topic: str | None = None,
state_path: str | None = None,
+ state_on_value: str | None = None,
+ # Legacy: single topic/path/multiplier (for backward compatibility)
+ topic: str | None = None,
multiplier: float = 1.0,
):
- """Subscribe to a topic for a plug."""
+ """Subscribe to MQTT topics for a plug.
+
+ Each data type (power, energy, state) can have its own topic.
+ For backward compatibility, if power_topic is not set but topic is,
+ topic will be used for all data types that have paths configured.
+ """
with self._lock:
- # Store configuration
- self.plug_configs[plug_id] = (topic, power_path, energy_path, state_path, multiplier)
+ # Initialize config for this plug
+ self.plug_configs[plug_id] = {}
- # Add to subscriptions
- if topic not in self.subscriptions:
- self.subscriptions[topic] = []
- # Actually subscribe if connected
- if self.client and self.connected:
- try:
- self.client.subscribe(topic, qos=1)
- logger.info(f"MQTT smart plug {plug_id}: subscribed to {topic}")
- except Exception as e:
- logger.error(f"MQTT smart plug: failed to subscribe to {topic}: {e}")
+ # Determine topics (new fields take priority, fall back to legacy)
+ effective_power_topic = power_topic or topic
+ effective_energy_topic = energy_topic or topic
+ effective_state_topic = state_topic or topic
- if plug_id not in self.subscriptions[topic]:
- self.subscriptions[topic].append(plug_id)
+ # Use new multipliers or fall back to legacy
+ effective_power_mult = power_multiplier if power_multiplier != 1.0 else multiplier
+ effective_energy_mult = energy_multiplier if energy_multiplier != 1.0 else multiplier
+
+ # Configure power subscription
+ if effective_power_topic and power_path:
+ config = MQTTDataSourceConfig(
+ topic=effective_power_topic,
+ path=power_path,
+ multiplier=effective_power_mult,
+ )
+ self.plug_configs[plug_id]["power"] = config
+ self._add_subscription(plug_id, effective_power_topic, "power")
+
+ # Configure energy subscription
+ if effective_energy_topic and energy_path:
+ config = MQTTDataSourceConfig(
+ topic=effective_energy_topic,
+ path=energy_path,
+ multiplier=effective_energy_mult,
+ )
+ self.plug_configs[plug_id]["energy"] = config
+ self._add_subscription(plug_id, effective_energy_topic, "energy")
+
+ # Configure state subscription
+ if effective_state_topic and state_path:
+ config = MQTTDataSourceConfig(
+ topic=effective_state_topic,
+ path=state_path,
+ on_value=state_on_value,
+ )
+ self.plug_configs[plug_id]["state"] = config
+ self._add_subscription(plug_id, effective_state_topic, "state")
# Initialize data entry
if plug_id not in self.plug_data:
self.plug_data[plug_id] = SmartPlugMQTTData(plug_id=plug_id)
+ logger.info(
+ f"MQTT smart plug {plug_id}: configured with "
+ f"power={effective_power_topic if power_path else None}, "
+ f"energy={effective_energy_topic if energy_path else None}, "
+ f"state={effective_state_topic if state_path else None}"
+ )
+
+ def _add_subscription(self, plug_id: int, topic: str, data_type: str):
+ """Add a subscription for a plug/data_type to a topic."""
+ if topic not in self.subscriptions:
+ self.subscriptions[topic] = []
+ # Actually subscribe if connected
+ if self.client and self.connected:
+ try:
+ self.client.subscribe(topic, qos=1)
+ logger.info(f"MQTT smart plug: subscribed to {topic}")
+ except Exception as e:
+ logger.error(f"MQTT smart plug: failed to subscribe to {topic}: {e}")
+
+ entry = (plug_id, data_type)
+ if entry not in self.subscriptions[topic]:
+ self.subscriptions[topic].append(entry)
+
def unsubscribe(self, plug_id: int):
"""Unsubscribe when plug is deleted/updated."""
with self._lock:
- # Get the topic for this plug
- config = self.plug_configs.pop(plug_id, None)
- if not config:
- return
+ # Get all configs for this plug
+ configs = self.plug_configs.pop(plug_id, {})
+ if not configs:
+ # Still clean up any stray subscriptions
+ pass
- topic = config[0]
+ # Collect all topics this plug was subscribed to
+ topics_to_check = set()
+ for _data_type, config in configs.items():
+ topics_to_check.add(config.topic)
- # Remove from subscriptions
- if topic in self.subscriptions:
- if plug_id in self.subscriptions[topic]:
- self.subscriptions[topic].remove(plug_id)
+ # Also scan subscriptions to remove any entries for this plug
+ for topic in list(self.subscriptions.keys()):
+ # Remove all entries for this plug_id
+ self.subscriptions[topic] = [(pid, dtype) for pid, dtype in self.subscriptions[topic] if pid != plug_id]
+ topics_to_check.add(topic)
- # If no more plugs on this topic, unsubscribe
- if not self.subscriptions[topic]:
+ # Unsubscribe from topics with no more subscribers
+ for topic in topics_to_check:
+ if topic in self.subscriptions and not self.subscriptions[topic]:
del self.subscriptions[topic]
if self.client and self.connected:
try:
diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py
index 69b520fd1..1b4bb262b 100644
--- a/backend/tests/conftest.py
+++ b/backend/tests/conftest.py
@@ -315,9 +315,19 @@ def smart_plug_factory(db_session):
defaults["ha_entity_id"] = "switch.test"
defaults["ip_address"] = None
elif plug_type == "mqtt":
+ # Legacy fields (for backward compatibility tests)
defaults["mqtt_topic"] = kwargs.get("mqtt_topic", "test/topic")
- defaults["mqtt_power_path"] = kwargs.get("mqtt_power_path", "power")
defaults["mqtt_multiplier"] = kwargs.get("mqtt_multiplier", 1.0)
+ # New separate topic/path/multiplier fields
+ defaults["mqtt_power_topic"] = kwargs.get("mqtt_power_topic")
+ defaults["mqtt_power_path"] = kwargs.get("mqtt_power_path", "power")
+ defaults["mqtt_power_multiplier"] = kwargs.get("mqtt_power_multiplier", 1.0)
+ defaults["mqtt_energy_topic"] = kwargs.get("mqtt_energy_topic")
+ defaults["mqtt_energy_path"] = kwargs.get("mqtt_energy_path")
+ defaults["mqtt_energy_multiplier"] = kwargs.get("mqtt_energy_multiplier", 1.0)
+ defaults["mqtt_state_topic"] = kwargs.get("mqtt_state_topic")
+ defaults["mqtt_state_path"] = kwargs.get("mqtt_state_path")
+ defaults["mqtt_state_on_value"] = kwargs.get("mqtt_state_on_value")
defaults["ip_address"] = None
defaults["ha_entity_id"] = None
else:
diff --git a/backend/tests/integration/test_smart_plugs_api.py b/backend/tests/integration/test_smart_plugs_api.py
index 90b7b1e8f..0a6421ae0 100644
--- a/backend/tests/integration/test_smart_plugs_api.py
+++ b/backend/tests/integration/test_smart_plugs_api.py
@@ -694,3 +694,132 @@ class TestSmartPlugsAPI:
result = response.json()
assert result["mqtt_topic"] == "new/topic"
assert result["mqtt_power_path"] == "new_power"
+
+ # ========================================================================
+ # Enhanced MQTT Integration tests (separate topics per data type)
+ # ========================================================================
+
+ @pytest.mark.asyncio
+ @pytest.mark.integration
+ async def test_create_mqtt_plug_with_separate_topics(self, async_client: AsyncClient, mock_mqtt_smart_plug_service):
+ """Verify MQTT plug can be created with separate topics for power, energy, and state."""
+ data = {
+ "name": "MQTT Separate Topics",
+ "plug_type": "mqtt",
+ "mqtt_power_topic": "zigbee/power",
+ "mqtt_power_path": "power_l1",
+ "mqtt_power_multiplier": 0.001,
+ "mqtt_energy_topic": "zigbee/energy",
+ "mqtt_energy_path": "energy_total",
+ "mqtt_energy_multiplier": 1.0,
+ "mqtt_state_topic": "zigbee/state",
+ "mqtt_state_path": "state",
+ "mqtt_state_on_value": "ON",
+ "enabled": True,
+ }
+
+ response = await async_client.post("/api/v1/smart-plugs/", json=data)
+
+ assert response.status_code == 200
+ result = response.json()
+ assert result["name"] == "MQTT Separate Topics"
+ assert result["plug_type"] == "mqtt"
+ # Power fields
+ assert result["mqtt_power_topic"] == "zigbee/power"
+ assert result["mqtt_power_path"] == "power_l1"
+ assert result["mqtt_power_multiplier"] == 0.001
+ # Energy fields
+ assert result["mqtt_energy_topic"] == "zigbee/energy"
+ assert result["mqtt_energy_path"] == "energy_total"
+ assert result["mqtt_energy_multiplier"] == 1.0
+ # State fields
+ assert result["mqtt_state_topic"] == "zigbee/state"
+ assert result["mqtt_state_path"] == "state"
+ assert result["mqtt_state_on_value"] == "ON"
+
+ @pytest.mark.asyncio
+ @pytest.mark.integration
+ async def test_create_mqtt_plug_energy_only(self, async_client: AsyncClient, mock_mqtt_smart_plug_service):
+ """Verify MQTT plug can be created with only energy monitoring."""
+ data = {
+ "name": "Energy Only Monitor",
+ "plug_type": "mqtt",
+ "mqtt_energy_topic": "sensors/energy",
+ "mqtt_energy_path": "kwh",
+ "mqtt_energy_multiplier": 0.001, # Wh to kWh
+ "enabled": True,
+ }
+
+ response = await async_client.post("/api/v1/smart-plugs/", json=data)
+
+ assert response.status_code == 200
+ result = response.json()
+ assert result["mqtt_energy_topic"] == "sensors/energy"
+ assert result["mqtt_energy_path"] == "kwh"
+ assert result["mqtt_energy_multiplier"] == 0.001
+
+ @pytest.mark.asyncio
+ @pytest.mark.integration
+ async def test_create_mqtt_plug_state_only(self, async_client: AsyncClient, mock_mqtt_smart_plug_service):
+ """Verify MQTT plug can be created with only state monitoring."""
+ data = {
+ "name": "State Only Monitor",
+ "plug_type": "mqtt",
+ "mqtt_state_topic": "switches/outlet",
+ "mqtt_state_path": "state",
+ "mqtt_state_on_value": "true",
+ "enabled": True,
+ }
+
+ response = await async_client.post("/api/v1/smart-plugs/", json=data)
+
+ assert response.status_code == 200
+ result = response.json()
+ assert result["mqtt_state_topic"] == "switches/outlet"
+ assert result["mqtt_state_path"] == "state"
+ assert result["mqtt_state_on_value"] == "true"
+
+ @pytest.mark.asyncio
+ @pytest.mark.integration
+ async def test_create_mqtt_plug_no_data_source_fails(self, async_client: AsyncClient):
+ """Verify creating MQTT plug without any complete data source fails."""
+ data = {
+ "name": "Invalid MQTT Plug",
+ "plug_type": "mqtt",
+ # Has topic but no path - not a complete data source
+ "mqtt_power_topic": "zigbee/power",
+ "enabled": True,
+ }
+
+ response = await async_client.post("/api/v1/smart-plugs/", json=data)
+
+ assert response.status_code == 422 # Validation error
+
+ @pytest.mark.asyncio
+ @pytest.mark.integration
+ async def test_update_mqtt_plug_separate_multipliers(
+ self, async_client: AsyncClient, smart_plug_factory, db_session, mock_mqtt_smart_plug_service
+ ):
+ """Verify MQTT plug multipliers can be updated separately."""
+ plug = await smart_plug_factory(
+ plug_type="mqtt",
+ mqtt_power_topic="test/power",
+ mqtt_power_path="power",
+ mqtt_power_multiplier=1.0,
+ mqtt_energy_topic="test/energy",
+ mqtt_energy_path="energy",
+ mqtt_energy_multiplier=1.0,
+ )
+
+ response = await async_client.patch(
+ f"/api/v1/smart-plugs/{plug.id}",
+ json={
+ "mqtt_power_multiplier": 0.001, # Change power multiplier only
+ "mqtt_energy_multiplier": 0.001, # Change energy multiplier only
+ },
+ )
+
+ assert response.status_code == 200
+ result = response.json()
+ assert result["mqtt_power_multiplier"] == 0.001
+ assert result["mqtt_energy_multiplier"] == 0.001
diff --git a/frontend/src/__tests__/components/SmartPlugCard.test.tsx b/frontend/src/__tests__/components/SmartPlugCard.test.tsx
index da294867b..5182b08cb 100644
--- a/frontend/src/__tests__/components/SmartPlugCard.test.tsx
+++ b/frontend/src/__tests__/components/SmartPlugCard.test.tsx
@@ -21,12 +21,24 @@ const createMockPlug = (overrides: Partial
The MQTT topic to subscribe to
+ {/* Power Section */} +Power Monitoring
++ Use multiplier 0.001 for mW→W, 1000 for kW→W +
Path to power value in JSON (dot notation)
+ {/* Energy Section */} +Energy Monitoring (optional)
++ Use multiplier 0.001 for Wh→kWh, 1000 for MWh→kWh +
Path to ON/OFF state (optional)
-Path to energy/kWh value (optional)
-- Multiply values by this factor. Use 0.001 to convert mW to W, 1000 for kW to W. + {/* State Section */} +
State Monitoring (optional)
++ ON value: the exact string that means "ON". Leave empty for auto-detect (ON, true, 1)
yn?(Dn=Yt,Yt=null):Dn=Yt.sibling;var Gn=st(Qe,Yt,nt[yn],vt);if(Gn===null){Yt===null&&(Yt=Dn);break}h&&Yt&&Gn.alternate===null&&p(Qe,Yt),Ge=L(Gn,Ge,yn),Vn===null?Zt=Gn:Vn.sibling=Gn,Vn=Gn,Yt=Dn}if(yn===nt.length)return _(Qe,Yt),Ln&&ac(Qe,yn),Zt;if(Yt===null){for(;yn