chore: add mqtt json getter to timeout handler (#13237)
This commit is contained in:
parent
412ec23a77
commit
681dcb069f
2 changed files with 44 additions and 70 deletions
|
|
@ -19,14 +19,13 @@ type Warp2 struct {
|
|||
log *util.Logger
|
||||
client *mqtt.Client
|
||||
features []string
|
||||
maxcurrentG func() (string, error)
|
||||
statusG func() (string, error)
|
||||
meterG func() (string, error)
|
||||
meterDetailsG func() (string, error)
|
||||
chargeG func() (string, error)
|
||||
userconfigG func() (string, error)
|
||||
emStateG func() (string, error)
|
||||
emLowLevelG func() (string, error)
|
||||
maxcurrentG func(any) error
|
||||
statusG func(any) error
|
||||
meterG func(any) error
|
||||
meterDetailsG func(any) error
|
||||
chargeG func(any) error
|
||||
emStateG func(any) error
|
||||
emLowLevelG func(any) error
|
||||
maxcurrentS func(int64) error
|
||||
phasesS func(int64) error
|
||||
current int64
|
||||
|
|
@ -115,27 +114,23 @@ func NewWarp2(mqttconf mqtt.Config, topic, emTopic string, timeout time.Duration
|
|||
return provider.NewMqtt(log, client, fmt.Sprintf(s, args...), 0)
|
||||
}
|
||||
|
||||
wb.maxcurrentG, err = to.StringGetter(mq("%s/evse/external_current", topic))
|
||||
wb.maxcurrentG, err = to.JsonGetter(mq("%s/evse/external_current", topic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
wb.statusG, err = to.StringGetter(mq("%s/evse/state", topic))
|
||||
wb.statusG, err = to.JsonGetter(mq("%s/evse/state", topic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
wb.meterG, err = to.StringGetter(mq("%s/meter/values", topic))
|
||||
wb.meterG, err = to.JsonGetter(mq("%s/meter/values", topic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
wb.meterDetailsG, err = to.StringGetter(mq("%s/meter/all_values", topic))
|
||||
wb.meterDetailsG, err = to.JsonGetter(mq("%s/meter/all_values", topic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
wb.chargeG, err = to.StringGetter(mq("%s/charge_tracker/current_charge", topic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
wb.userconfigG, err = to.StringGetter(mq("%s/users/config", topic))
|
||||
wb.chargeG, err = to.JsonGetter(mq("%s/charge_tracker/current_charge", topic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -148,7 +143,7 @@ func NewWarp2(mqttconf mqtt.Config, topic, emTopic string, timeout time.Duration
|
|||
return nil, err
|
||||
}
|
||||
|
||||
wb.emStateG, err = to.StringGetter(mq("%s/power_manager/state", emTopic))
|
||||
wb.emStateG, err = to.JsonGetter(mq("%s/power_manager/state", emTopic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -161,7 +156,7 @@ func NewWarp2(mqttconf mqtt.Config, topic, emTopic string, timeout time.Duration
|
|||
return nil, err
|
||||
}
|
||||
|
||||
wb.emLowLevelG, err = to.StringGetter(mq("%s/power_manager/low_level_state", emTopic))
|
||||
wb.emLowLevelG, err = to.JsonGetter(mq("%s/power_manager/low_level_state", emTopic))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -199,12 +194,7 @@ func (wb *Warp2) Enable(enable bool) error {
|
|||
// Enabled implements the api.Charger interface
|
||||
func (wb *Warp2) Enabled() (bool, error) {
|
||||
var res warp.EvseExternalCurrent
|
||||
|
||||
s, err := wb.maxcurrentG()
|
||||
if err == nil {
|
||||
err = json.Unmarshal([]byte(s), &res)
|
||||
}
|
||||
|
||||
err := wb.maxcurrentG(&res)
|
||||
return res.Current >= 6000, err
|
||||
}
|
||||
|
||||
|
|
@ -212,13 +202,9 @@ func (wb *Warp2) Enabled() (bool, error) {
|
|||
func (wb *Warp2) Status() (api.ChargeStatus, error) {
|
||||
res := api.StatusNone
|
||||
|
||||
s, err := wb.statusG()
|
||||
if err != nil {
|
||||
return res, err
|
||||
}
|
||||
|
||||
var status warp.EvseState
|
||||
if err := json.Unmarshal([]byte(s), &status); err != nil {
|
||||
err := wb.statusG(&status)
|
||||
if err != nil {
|
||||
return res, err
|
||||
}
|
||||
|
||||
|
|
@ -256,39 +242,22 @@ func (wb *Warp2) MaxCurrentMillis(current float64) error {
|
|||
// CurrentPower implements the api.Meter interface
|
||||
func (wb *Warp2) currentPower() (float64, error) {
|
||||
var res warp.MeterValues
|
||||
|
||||
s, err := wb.meterG()
|
||||
if err == nil {
|
||||
err = json.Unmarshal([]byte(s), &res)
|
||||
}
|
||||
|
||||
err := wb.meterG(&res)
|
||||
return res.Power, err
|
||||
}
|
||||
|
||||
// TotalEnergy implements the api.MeterEnergy interface
|
||||
func (wb *Warp2) totalEnergy() (float64, error) {
|
||||
var res warp.MeterValues
|
||||
|
||||
s, err := wb.meterG()
|
||||
if err == nil {
|
||||
err = json.Unmarshal([]byte(s), &res)
|
||||
}
|
||||
|
||||
err := wb.meterG(&res)
|
||||
return res.EnergyAbs, err
|
||||
}
|
||||
|
||||
func (wb *Warp2) meterValues() ([]float64, error) {
|
||||
s, err := wb.meterDetailsG()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var res []float64
|
||||
if err := json.Unmarshal([]byte(s), &res); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
err := wb.meterDetailsG(&res)
|
||||
|
||||
if len(res) <= 5 {
|
||||
if err == nil && len(res) <= 5 {
|
||||
return nil, errors.New("invalid length")
|
||||
}
|
||||
|
||||
|
|
@ -317,34 +286,19 @@ func (wb *Warp2) voltages() (float64, float64, float64, error) {
|
|||
|
||||
func (wb *Warp2) identify() (string, error) {
|
||||
var res warp.ChargeTrackerCurrentCharge
|
||||
|
||||
s, err := wb.chargeG()
|
||||
if err == nil {
|
||||
err = json.Unmarshal([]byte(s), &res)
|
||||
}
|
||||
|
||||
err := wb.chargeG(&res)
|
||||
return res.AuthorizationInfo.TagId, err
|
||||
}
|
||||
|
||||
func (wb *Warp2) emState() (warp.EmState, error) {
|
||||
var res warp.EmState
|
||||
|
||||
s, err := wb.emStateG()
|
||||
if err == nil {
|
||||
err = json.Unmarshal([]byte(s), &res)
|
||||
}
|
||||
|
||||
err := wb.emStateG(&res)
|
||||
return res, err
|
||||
}
|
||||
|
||||
func (wb *Warp2) emLowLevelState() (warp.EmLowLevelState, error) {
|
||||
var res warp.EmLowLevelState
|
||||
|
||||
s, err := wb.emLowLevelG()
|
||||
if err == nil {
|
||||
err = json.Unmarshal([]byte(s), &res)
|
||||
}
|
||||
|
||||
err := wb.emLowLevelG(&res)
|
||||
return res, err
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,5 +1,7 @@
|
|||
package provider
|
||||
|
||||
import "encoding/json"
|
||||
|
||||
// TimeoutHandler is a wrapper for a Getter that times out after a given duration
|
||||
type TimeoutHandler struct {
|
||||
ticker func() (string, error)
|
||||
|
|
@ -50,3 +52,21 @@ func (h *TimeoutHandler) StringGetter(p StringProvider) (func() (string, error),
|
|||
return val, err
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (h *TimeoutHandler) JsonGetter(p StringProvider) (func(any) error, error) {
|
||||
g, err := p.StringGetter()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return func(res any) error {
|
||||
val, err := g()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := h.ticker(); err != nil {
|
||||
return err
|
||||
}
|
||||
return json.Unmarshal([]byte(val), res)
|
||||
}, nil
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue