Skip to content

Commit 45b2273

Browse files
authored
Merge pull request #22 from rddl-network/eckelj/single_values_mqtt_topic
added device status mqtt support and disabled the old aggregated powe…
2 parents 9984018 + 4d12a5b commit 45b2273

4 files changed

Lines changed: 94 additions & 5 deletions

File tree

README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,7 @@ topic = "energy-consumption-reports"
8787
{"value": 50.000, "timestamp": "2025-07-15 00:15:00"},
8888
{"value": 50.100, "timestamp": "2025-07-15 00:30:00"},
8989
... (total 96 entries) ...
90+
{"value": 60.100, "timestamp": "2025-07-16 00:00:00"}, // <-- the last entry is always 00:00:00 of the next day
9091
]
9192
}
9293
```

internal/model/devicestatus.go

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
package model
2+
3+
type DeviceStatus struct {
4+
IsOn bool `json:"isOn"`
5+
CurrentVoltage float64 `json:"currentVoltage"`
6+
CurrentAmps float64 `json:"currentAmps"`
7+
CurrentActivePower float64 `json:"currentActivePower"`
8+
TotalEnergyConsumed *float64 `json:"totalEnergyConsumed"`
9+
}
10+
11+
type DeviceStatusExt struct {
12+
ID string `json:"id"`
13+
DeviceStatus DeviceStatus `json:"deviceStatus"`
14+
}

internal/server/helper.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,31 @@ func (s *Server) write2InfluxDB(data model.EnergyData) error {
5757
return nil
5858
}
5959

60+
func (s *Server) writeDeviceStatus2InfluxDB(data model.DeviceStatusExt) error {
61+
writeAPI := s.influxDBClient
62+
if writeAPI == nil {
63+
log.Printf("No InfluxDB write API set")
64+
return nil
65+
}
66+
67+
err := writeAPI.WritePoint(
68+
context.Background(),
69+
"device_status",
70+
map[string]string{
71+
"ID": data.ID,
72+
},
73+
map[string]interface{}{
74+
"kW/h": *data.DeviceStatus.TotalEnergyConsumed,
75+
},
76+
time.Now(),
77+
)
78+
if err != nil {
79+
log.Printf("Failed to write to InfluxDB: %v", err)
80+
return err
81+
}
82+
return nil
83+
}
84+
6085
// sendJSONResponse sends a JSON response with the given status code
6186
func sendJSONResponse(w http.ResponseWriter, resp Response, statusCode int) {
6287
w.Header().Set("Content-Type", "application/json")

internal/server/mqtt.go

Lines changed: 54 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import (
66
"encoding/json"
77
"fmt"
88
"log"
9+
"strings"
910

1011
mqtt "github.com/eclipse/paho.mqtt.golang"
1112
"github.com/rddl-network/energy-service/internal/config"
@@ -25,23 +26,71 @@ func (s *Server) initMQTT() {
2526
tlsConfig := &tls.Config{}
2627
opts.SetTLSConfig(tlsConfig)
2728

29+
// Add connection lost handler
30+
opts.SetConnectionLostHandler(func(client mqtt.Client, err error) {
31+
log.Printf("MQTT connection lost: %v", err)
32+
})
33+
opts.AutoReconnect = true
34+
2835
client := mqtt.NewClient(opts)
2936
if token := client.Connect(); token.Wait() && token.Error() != nil {
3037
log.Printf("MQTT connect error: %v", token.Error())
3138
return
3239
}
3340
s.mqttClient = client
34-
topic := mqttCfg.Topic
35-
if topic == "" {
36-
topic = "energy-consumption-reports"
37-
}
38-
if token := client.Subscribe(topic, 0, s.handleMQTTMessage); token.Wait() && token.Error() != nil {
41+
//subscribe(client, mqttCfg.Topic, s.handleMQTTMessage)
42+
subscribe(client, "dirigera/+", s.handleSimpleDataMQTTMessage)
43+
}
44+
45+
func subscribe(client mqtt.Client, topic string, callback mqtt.MessageHandler) {
46+
if token := client.Subscribe(topic, 0, callback); token.Wait() && token.Error() != nil {
3947
log.Printf("MQTT subscribe error: %v", token.Error())
4048
} else {
4149
log.Printf("Subscribed to MQTT topic: %s", topic)
4250
}
4351
}
4452

53+
func extractIDFromTopic(topic string) string {
54+
parts := strings.Split(topic, "/")
55+
if len(parts) == 2 && parts[0] == "dirigera" {
56+
return parts[1]
57+
}
58+
return ""
59+
}
60+
61+
func (s *Server) handleSimpleDataMQTTMessage(client mqtt.Client, msg mqtt.Message) {
62+
defer func() {
63+
if r := recover(); r != nil {
64+
log.Printf("MQTT handler panic: %v", r)
65+
}
66+
}()
67+
68+
var deviceStatusExt model.DeviceStatusExt
69+
70+
deviceStatusExt.ID = extractIDFromTopic(msg.Topic())
71+
if deviceStatusExt.ID == "" {
72+
log.Printf("MQTT: Invalid topic format: %s", msg.Topic())
73+
return
74+
}
75+
76+
if err := json.Unmarshal(msg.Payload(), &deviceStatusExt.DeviceStatus); err != nil {
77+
log.Printf("MQTT: Failed to decode JSON: %v", err)
78+
return
79+
}
80+
81+
if deviceStatusExt.DeviceStatus.TotalEnergyConsumed == nil {
82+
log.Printf("MQTT: JSON does not contain consumed energy values.")
83+
return
84+
}
85+
86+
err := s.writeDeviceStatus2InfluxDB(deviceStatusExt)
87+
if err != nil {
88+
log.Printf("MQTT: Failed to write to database: %v", err)
89+
return
90+
}
91+
log.Printf("MQTT: Energy data received and written to database successfully")
92+
}
93+
4594
// handleMQTTMessage processes incoming MQTT messages as energy data
4695
func (s *Server) handleMQTTMessage(client mqtt.Client, msg mqtt.Message) {
4796
var energyData model.EnergyData

0 commit comments

Comments
 (0)