Improve mqtt publishing (#421)
This commit is contained in:
parent
b11381b7e6
commit
9090514a9e
11 changed files with 100 additions and 43 deletions
|
|
@ -464,11 +464,11 @@ The MQTT API follows the REST API's structure:
|
|||
- `evcc/updated`: timestamp of last update
|
||||
- `evcc/site`: site dynamic state
|
||||
- `evcc/site/mode`: global charge mode (writable)
|
||||
- `evcc/site/targetsoc`: global target SoC (writable)
|
||||
- `evcc/site/targetSoC`: global target SoC (writable)
|
||||
- `evcc/loadpoints`: number of available loadpoints
|
||||
- `evcc/loadpoints/<id>`: loadpoint dynamic state
|
||||
- `evcc/loadpoints/<id>/mode`: loadpoint charge mode (writable)
|
||||
- `evcc/loadpoints/<id>/targetsoc`: loadpoint target SoC (writable)
|
||||
- `evcc/loadpoints/<id>/targetSoC`: loadpoint target SoC (writable)
|
||||
|
||||
Note: to modify writable settings append `/set` to the topic for writing.
|
||||
|
||||
|
|
|
|||
|
|
@ -19,7 +19,7 @@ import (
|
|||
|
||||
const (
|
||||
udpTimeout = time.Second
|
||||
kebaPort = "7090"
|
||||
kebaPort = 7090
|
||||
)
|
||||
|
||||
// RFID contains access credentials
|
||||
|
|
@ -63,16 +63,14 @@ func NewKeba(conn, serial string, rfid RFID, timeout time.Duration) (api.Charger
|
|||
|
||||
var err error
|
||||
if keba.Instance == nil {
|
||||
keba.Instance, err = keba.New(log, fmt.Sprintf(":%s", kebaPort))
|
||||
keba.Instance, err = keba.New(log, fmt.Sprintf(":%d", kebaPort))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
// add default port
|
||||
if _, _, err = net.SplitHostPort(conn); err != nil {
|
||||
conn = fmt.Sprintf("%s:%s", conn, kebaPort)
|
||||
}
|
||||
conn = util.DefaultPort(conn, kebaPort)
|
||||
|
||||
c := &Keba{
|
||||
log: log,
|
||||
|
|
|
|||
|
|
@ -3,7 +3,6 @@ package charger
|
|||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"net"
|
||||
"time"
|
||||
|
||||
"github.com/andig/evcc/api"
|
||||
|
|
@ -52,10 +51,6 @@ func NewPhoenixEMCPFromConfig(other map[string]interface{}) (api.Charger, error)
|
|||
return nil, err
|
||||
}
|
||||
|
||||
if _, _, err := net.SplitHostPort(cc.URI); err != nil {
|
||||
return nil, fmt.Errorf("missing or invalid phoenix uri: %s", cc.URI)
|
||||
}
|
||||
|
||||
wb, err := NewPhoenixEMCP(cc.URI, cc.ID)
|
||||
|
||||
var currentPower func() (float64, error)
|
||||
|
|
|
|||
|
|
@ -147,7 +147,7 @@ func run(cmd *cobra.Command, args []string) {
|
|||
}
|
||||
|
||||
// setup mqtt publisher
|
||||
if conf.Mqtt.Broker != "" && conf.Mqtt.Topic != "" {
|
||||
if conf.Mqtt.Broker != "" {
|
||||
publisher := server.NewMQTT(conf.Mqtt.Topic)
|
||||
go publisher.Run(site, tee.Attach())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -342,6 +342,12 @@ func (lp *LoadPoint) Prepare(uiChan chan<- util.Param, pushChan chan<- push.Even
|
|||
_ = lp.bus.Subscribe(evChargeCurrent, lp.evChargeCurrentHandler)
|
||||
|
||||
// publish initial values
|
||||
lp.publish("title", lp.Title)
|
||||
lp.publish("minCurrent", lp.MinCurrent)
|
||||
lp.publish("maxCurrent", lp.MaxCurrent)
|
||||
lp.publish("phases", lp.Phases)
|
||||
lp.publish("activePhases", lp.Phases)
|
||||
|
||||
lp.Lock()
|
||||
lp.publish("mode", lp.Mode)
|
||||
lp.publish("targetSoC", lp.SoC.Target)
|
||||
|
|
@ -411,8 +417,8 @@ func (lp *LoadPoint) setActiveVehicle(vehicle api.Vehicle) {
|
|||
lp.vehicle = vehicle
|
||||
lp.socEstimator = wrapper.NewSocEstimator(lp.log, vehicle, lp.SoC.Estimate)
|
||||
|
||||
lp.publish("socCapacity", lp.vehicle.Capacity())
|
||||
lp.publish("socTitle", lp.vehicle.Title())
|
||||
lp.publish("socCapacity", lp.vehicle.Capacity())
|
||||
}
|
||||
|
||||
// findActiveVehicle validates if the active vehicle is still connected to the loadpoint
|
||||
|
|
|
|||
81
core/site.go
81
core/site.go
|
|
@ -158,46 +158,85 @@ func (site *Site) Configuration() SiteConfiguration {
|
|||
return c
|
||||
}
|
||||
|
||||
func logMeter(log *util.Logger, meter interface{}) {
|
||||
func meterCapabilities(name string, meter interface{}) string {
|
||||
_, power := meter.(api.Meter)
|
||||
_, energy := meter.(api.MeterEnergy)
|
||||
_, currents := meter.(api.MeterCurrent)
|
||||
|
||||
log.INFO.Printf(" power %s", presence[power])
|
||||
log.INFO.Printf(" energy %s", presence[energy])
|
||||
log.INFO.Printf(" currents %s", presence[currents])
|
||||
name += ":"
|
||||
return fmt.Sprintf(" %-8s power %s energy %s currents %s",
|
||||
name,
|
||||
presence[power],
|
||||
presence[energy],
|
||||
presence[currents],
|
||||
)
|
||||
}
|
||||
|
||||
// DumpConfig site configuration
|
||||
func (site *Site) DumpConfig() {
|
||||
site.log.INFO.Println("site config:")
|
||||
site.log.INFO.Printf(" grid %s", presence[site.gridMeter != nil])
|
||||
site.log.INFO.Printf(" pv %s", presence[site.pvMeter != nil])
|
||||
site.log.INFO.Printf(" battery %s", presence[site.batteryMeter != nil])
|
||||
site.publish("title", site.Title)
|
||||
|
||||
site.log.INFO.Println("site config:")
|
||||
site.log.INFO.Printf(" meters: grid %s pv %s battery %s",
|
||||
presence[site.gridMeter != nil],
|
||||
presence[site.pvMeter != nil],
|
||||
presence[site.batteryMeter != nil],
|
||||
)
|
||||
|
||||
site.publish("gridConfigured", site.gridMeter != nil)
|
||||
if site.gridMeter != nil {
|
||||
site.log.INFO.Println(" grid meter config:")
|
||||
logMeter(site.log, site.gridMeter)
|
||||
site.log.INFO.Println(meterCapabilities("grid", site.gridMeter))
|
||||
}
|
||||
|
||||
site.publish("pvConfigured", site.pvMeter != nil)
|
||||
if site.pvMeter != nil {
|
||||
site.log.INFO.Println(meterCapabilities("pv", site.pvMeter))
|
||||
}
|
||||
|
||||
site.publish("batteryConfigured", site.batteryMeter != nil)
|
||||
if site.batteryMeter != nil {
|
||||
_, ok := site.batteryMeter.(api.Battery)
|
||||
site.log.INFO.Println(
|
||||
meterCapabilities("battery", site.batteryMeter),
|
||||
fmt.Sprintf("soc %s", presence[ok]),
|
||||
)
|
||||
}
|
||||
|
||||
for i, lp := range site.loadpoints {
|
||||
lp.log.INFO.Printf("loadpoint %d config:", i+1)
|
||||
lp.log.INFO.Printf("loadpoint %d:", i+1)
|
||||
|
||||
lp.log.INFO.Printf(" vehicle %s", presence[lp.vehicle != nil])
|
||||
lp.log.INFO.Printf(" charge %s", presence[lp.HasChargeMeter()])
|
||||
if lp.HasChargeMeter() {
|
||||
lp.log.INFO.Println(" charge meter config:")
|
||||
logMeter(site.log, lp.chargeMeter)
|
||||
}
|
||||
lp.log.INFO.Printf(" mode: %s", lp.GetMode())
|
||||
|
||||
charger := lp.handler.(*ChargerHandler).charger
|
||||
_, power := charger.(api.Meter)
|
||||
_, energy := charger.(api.MeterEnergy)
|
||||
_, currents := charger.(api.MeterCurrent)
|
||||
_, timer := charger.(api.ChargeTimer)
|
||||
|
||||
lp.log.INFO.Println(" charger config:")
|
||||
logMeter(lp.log, charger)
|
||||
lp.log.INFO.Printf(" timer %s", presence[timer])
|
||||
lp.log.INFO.Printf(" charger: power %s energy %s currents %s timer %s",
|
||||
presence[power],
|
||||
presence[energy],
|
||||
presence[currents],
|
||||
presence[timer],
|
||||
)
|
||||
|
||||
lp.log.INFO.Printf(" mode: %s", lp.GetMode())
|
||||
lp.log.INFO.Printf(" meters: charge %s", presence[lp.HasChargeMeter()])
|
||||
|
||||
lp.publish("chargeConfigured", lp.HasChargeMeter())
|
||||
if lp.HasChargeMeter() {
|
||||
lp.log.INFO.Printf(meterCapabilities(" charge", lp.chargeMeter))
|
||||
}
|
||||
|
||||
lp.log.INFO.Printf(" vehicles: %s", presence[len(lp.vehicles) > 0])
|
||||
|
||||
for i, v := range lp.vehicles {
|
||||
_, estimate := v.(api.ChargeFinishTimer)
|
||||
_, status := v.(api.VehicleStatus)
|
||||
_, climate := v.(api.Climater)
|
||||
lp.log.INFO.Printf(" car %d: estimate %s status %s climate %s",
|
||||
i, presence[estimate], presence[status], presence[climate],
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -47,6 +47,8 @@ func NewMqttClient(
|
|||
qos byte,
|
||||
) *MqttClient {
|
||||
log := util.NewLogger("mqtt")
|
||||
|
||||
broker = util.DefaultPort(broker, 1883)
|
||||
log.INFO.Printf("connecting %s at %s", clientID, broker)
|
||||
|
||||
mc := &MqttClient{
|
||||
|
|
|
|||
|
|
@ -3,7 +3,6 @@ package server
|
|||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/andig/evcc/api"
|
||||
|
|
@ -20,6 +19,10 @@ type MQTT struct {
|
|||
|
||||
// NewMQTT creates MQTT server
|
||||
func NewMQTT(root string) *MQTT {
|
||||
if root == "" {
|
||||
root = "evcc"
|
||||
}
|
||||
|
||||
return &MQTT{
|
||||
Handler: provider.MQTT,
|
||||
root: root,
|
||||
|
|
@ -37,7 +40,7 @@ func (m *MQTT) encode(v interface{}) string {
|
|||
case fmt.Stringer, string:
|
||||
s = fmt.Sprintf("%s", val)
|
||||
case float64:
|
||||
s = fmt.Sprintf("%.3f", val)
|
||||
s = fmt.Sprintf("%.4g", val)
|
||||
default:
|
||||
s = fmt.Sprintf("%v", val)
|
||||
}
|
||||
|
|
@ -87,7 +90,7 @@ func (m *MQTT) Run(site core.SiteAPI, in <-chan util.Param) {
|
|||
m.publish(topic, true, len(site.LoadPoints()))
|
||||
|
||||
for id, lp := range site.LoadPoints() {
|
||||
topic := fmt.Sprintf("%s/loadpoints/%d", m.root, id)
|
||||
topic := fmt.Sprintf("%s/loadpoints/%d", m.root, id+1)
|
||||
m.listenSetters(topic, lp)
|
||||
}
|
||||
|
||||
|
|
@ -98,7 +101,8 @@ func (m *MQTT) Run(site core.SiteAPI, in <-chan util.Param) {
|
|||
for p := range in {
|
||||
topic := fmt.Sprintf("%s/site", m.root)
|
||||
if p.LoadPoint != nil {
|
||||
topic = fmt.Sprintf("%s/loadpoints/%d", m.root, *p.LoadPoint)
|
||||
id := *p.LoadPoint + 1
|
||||
topic = fmt.Sprintf("%s/loadpoints/%d", m.root, id)
|
||||
}
|
||||
|
||||
// alive indicator
|
||||
|
|
@ -108,7 +112,7 @@ func (m *MQTT) Run(site core.SiteAPI, in <-chan util.Param) {
|
|||
}
|
||||
|
||||
// value
|
||||
topic += "/" + strings.ToLower(p.Key)
|
||||
topic += "/" + p.Key
|
||||
m.publish(topic, false, p.Val)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -39,8 +39,7 @@ func Run(log *util.Logger, cache chan<- util.Param) {
|
|||
cache: cache,
|
||||
}
|
||||
|
||||
instance.checkVersion()
|
||||
for range time.NewTicker(24 * time.Hour).C {
|
||||
for range time.NewTicker(time.Hour).C {
|
||||
instance.checkVersion()
|
||||
}
|
||||
}
|
||||
|
|
|
|||
15
util/net.go
Normal file
15
util/net.go
Normal file
|
|
@ -0,0 +1,15 @@
|
|||
package util
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net"
|
||||
)
|
||||
|
||||
// DefaultPort appends given port to connection if not specified
|
||||
func DefaultPort(conn string, port int) string {
|
||||
if _, _, err := net.SplitHostPort(conn); err != nil {
|
||||
conn = fmt.Sprintf("%s:%d", conn, port)
|
||||
}
|
||||
|
||||
return conn
|
||||
}
|
||||
|
|
@ -71,7 +71,6 @@ func NewConfigurableFromConfig(other map[string]interface{}) (api.Vehicle, error
|
|||
chargeG: getter,
|
||||
}
|
||||
|
||||
// decorate vehicle with BatterySoC
|
||||
// decorate vehicle with Status
|
||||
var status func() (api.ChargeStatus, error)
|
||||
if cc.Status != nil {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue