diff --git a/provider/cache.go b/provider/cache.go index 1e8673944..026420766 100644 --- a/provider/cache.go +++ b/provider/cache.go @@ -3,11 +3,13 @@ package provider import ( "time" + "github.com/andig/evcc/api" "github.com/benbjohnson/clock" ) // Cached wraps a getter with a cache type Cached struct { + log *api.Logger clock clock.Clock updated time.Time cache time.Duration @@ -16,8 +18,9 @@ type Cached struct { } // NewCached wraps a getter with a cache -func NewCached(getter interface{}, cache time.Duration) *Cached { +func NewCached(log *api.Logger, getter interface{}, cache time.Duration) *Cached { return &Cached{ + log: log, clock: clock.New(), getter: getter, cache: cache, @@ -29,7 +32,7 @@ func (c *Cached) FloatGetter() FloatGetter { g, ok := c.getter.(FloatGetter) if !ok { if g, ok = c.getter.(func() (float64, error)); !ok { - log.FATAL.Fatalf("invalid type: %T", c.getter) + c.log.FATAL.Fatalf("invalid type: %T", c.getter) } g = FloatGetter(g) } @@ -54,7 +57,7 @@ func (c *Cached) IntGetter() IntGetter { g, ok := c.getter.(IntGetter) if !ok { if g, ok = c.getter.(func() (int64, error)); !ok { - log.FATAL.Fatalf("invalid type: %T", c.getter) + c.log.FATAL.Fatalf("invalid type: %T", c.getter) } g = IntGetter(g) } @@ -79,7 +82,7 @@ func (c *Cached) StringGetter() StringGetter { g, ok := c.getter.(StringGetter) if !ok { if g, ok = c.getter.(func() (string, error)); !ok { - log.FATAL.Fatalf("invalid type: %T", c.getter) + c.log.FATAL.Fatalf("invalid type: %T", c.getter) } g = StringGetter(g) } @@ -105,7 +108,7 @@ func (c *Cached) BoolGetter() BoolGetter { g, ok := c.getter.(BoolGetter) if !ok { if g, ok = c.getter.(func() (bool, error)); !ok { - log.FATAL.Fatalf("invalid type: %T", g) + c.log.FATAL.Fatalf("invalid type: %T", g) } g = BoolGetter(g) } diff --git a/provider/cache_test.go b/provider/cache_test.go index c0683eb8f..4eab785bc 100644 --- a/provider/cache_test.go +++ b/provider/cache_test.go @@ -27,7 +27,7 @@ func TestCachedGetter(t *testing.T) { } duration := time.Second - c := NewCached(g, duration) + c := NewCached(nil, g, duration) clck := clock.NewMock() c.clock = clck getter := c.FloatGetter() diff --git a/provider/config.go b/provider/config.go index a49a8c387..b5685039b 100644 --- a/provider/config.go +++ b/provider/config.go @@ -70,7 +70,7 @@ func NewFloatGetterFromConfig(log *api.Logger, config Config) (res FloatGetter) pc := scriptFromConfig(log, config.Other) res = NewScriptProvider(pc.Timeout).FloatGetter(pc.Cmd) if pc.Cache > 0 { - res = NewCached(res, pc.Cache).FloatGetter() + res = NewCached(log, res, pc.Cache).FloatGetter() } case "modbus-rtu", "modbus-tcp", "modbus-rtuovertcp", "modbus-tcprtu", "modbus-rtutcp": res = FloatGetter(NewModbusFromConfig(log, config.Type, config.Other).FloatGetter) @@ -91,7 +91,7 @@ func NewIntGetterFromConfig(log *api.Logger, config Config) (res IntGetter) { pc := scriptFromConfig(log, config.Other) res = NewScriptProvider(pc.Timeout).IntGetter(pc.Cmd) if pc.Cache > 0 { - res = NewCached(res, pc.Cache).IntGetter() + res = NewCached(log, res, pc.Cache).IntGetter() } case "modbus-rtu", "modbus-tcp", "modbus-rtuovertcp", "modbus-tcprtu", "modbus-rtutcp": res = IntGetter(NewModbusFromConfig(log, config.Type, config.Other).IntGetter) @@ -112,7 +112,7 @@ func NewStringGetterFromConfig(log *api.Logger, config Config) (res StringGetter pc := scriptFromConfig(log, config.Other) res = NewScriptProvider(pc.Timeout).StringGetter(pc.Cmd) if pc.Cache > 0 { - res = NewCached(res, pc.Cache).StringGetter() + res = NewCached(log, res, pc.Cache).StringGetter() } case "combined", "openwb": res = openWBStatusFromConfig(log, config.Other) @@ -133,7 +133,7 @@ func NewBoolGetterFromConfig(log *api.Logger, config Config) (res BoolGetter) { pc := scriptFromConfig(log, config.Other) res = NewScriptProvider(pc.Timeout).BoolGetter(pc.Cmd) if pc.Cache > 0 { - res = NewCached(res, pc.Cache).BoolGetter() + res = NewCached(log, res, pc.Cache).BoolGetter() } default: log.FATAL.Fatalf("invalid provider type %s", config.Type) diff --git a/provider/exec.go b/provider/exec.go index 311688ce9..15a47d104 100644 --- a/provider/exec.go +++ b/provider/exec.go @@ -12,10 +12,9 @@ import ( "github.com/kballard/go-shellquote" ) -var log = api.NewLogger("exec") - // Script implements shell script-based providers and setters type Script struct { + log *api.Logger timeout time.Duration } @@ -23,6 +22,7 @@ type Script struct { // Script execution is aborted after given timeout. func NewScriptProvider(timeout time.Duration) *Script { return &Script{ + log: api.NewLogger("exec"), timeout: timeout, } } @@ -53,11 +53,11 @@ func (e *Script) StringGetter(script string) StringGetter { s = strings.TrimSpace(string(ee.Stderr)) } - log.ERROR.Printf("%s: %s", strings.Join(args, " "), s) + e.log.ERROR.Printf("%s: %s", strings.Join(args, " "), s) return "", err } - log.TRACE.Printf("%s: %s", strings.Join(args, " "), s) + e.log.TRACE.Printf("%s: %s", strings.Join(args, " "), s) return s, nil } } diff --git a/provider/mqtt.go b/provider/mqtt.go index dc2be0c7a..5e0ec2c58 100644 --- a/provider/mqtt.go +++ b/provider/mqtt.go @@ -2,6 +2,7 @@ package provider import ( "fmt" + "math" "strconv" "sync" "time" @@ -15,10 +16,9 @@ const ( waitTimeout = 50 * time.Millisecond // polling interval when waiting for initial value ) -var mlog = api.NewLogger("mqtt") - // MqttClient is a paho publisher type MqttClient struct { + log *api.Logger mux sync.Mutex Client mqtt.Client broker string @@ -34,9 +34,11 @@ func NewMqttClient( clientID string, qos byte, ) *MqttClient { - mlog.INFO.Printf("connecting %s at %s", clientID, broker) + log := api.NewLogger("mqtt") + log.INFO.Printf("connecting %s at %s", clientID, broker) mc := &MqttClient{ + log: log, broker: broker, qos: qos, listener: make(map[string]func(string)), @@ -63,18 +65,18 @@ func NewMqttClient( // ConnectionLostHandler logs cause of connection loss as warning func (m *MqttClient) ConnectionLostHandler(client mqtt.Client, reason error) { - mlog.WARN.Printf("%s connection lost: %v", m.broker, reason.Error()) + m.log.WARN.Printf("%s connection lost: %v", m.broker, reason.Error()) } // ConnectionHandler restores listeners func (m *MqttClient) ConnectionHandler(client mqtt.Client) { - mlog.TRACE.Printf("%s connected", m.broker) + m.log.TRACE.Printf("%s connected", m.broker) m.mux.Lock() defer m.mux.Unlock() for topic, l := range m.listener { - mlog.TRACE.Printf("%s subscribe %s", m.broker, topic) + m.log.TRACE.Printf("%s subscribe %s", m.broker, topic) go m.listen(topic, l) } } @@ -83,7 +85,7 @@ func (m *MqttClient) ConnectionHandler(client mqtt.Client) { func (m *MqttClient) Listen(topic string, callback func(string)) { m.mux.Lock() if _, ok := m.listener[topic]; ok { - mlog.FATAL.Fatalf("%s: duplicate listener not allowed", topic) + m.log.FATAL.Fatalf("%s: duplicate listener not allowed", topic) } m.listener[topic] = callback m.mux.Unlock() @@ -105,6 +107,7 @@ func (m *MqttClient) listen(topic string, callback func(string)) { // FloatGetter creates handler for float64 from MQTT topic that returns cached value func (m *MqttClient) FloatGetter(topic string, multiplier float64, timeout time.Duration) FloatGetter { h := &msgHandler{ + log: m.log, topic: topic, multiplier: multiplier, timeout: timeout, @@ -117,6 +120,7 @@ func (m *MqttClient) FloatGetter(topic string, multiplier float64, timeout time. // IntGetter creates handler for int64 from MQTT topic that returns cached value func (m *MqttClient) IntGetter(topic string, multiplier int64, timeout time.Duration) IntGetter { h := &msgHandler{ + log: m.log, topic: topic, multiplier: float64(multiplier), timeout: timeout, @@ -129,6 +133,7 @@ func (m *MqttClient) IntGetter(topic string, multiplier int64, timeout time.Dura // StringGetter creates handler for string from MQTT topic that returns cached value func (m *MqttClient) StringGetter(topic string, timeout time.Duration) StringGetter { h := &msgHandler{ + log: m.log, topic: topic, timeout: timeout, } @@ -140,6 +145,7 @@ func (m *MqttClient) StringGetter(topic string, timeout time.Duration) StringGet // BoolGetter creates handler for string from MQTT topic that returns cached value func (m *MqttClient) BoolGetter(topic string, timeout time.Duration) BoolGetter { h := &msgHandler{ + log: m.log, topic: topic, timeout: timeout, } @@ -167,7 +173,7 @@ func (m *MqttClient) IntSetter(param, topic, message string) IntSetter { return err } - mlog.TRACE.Printf("send %s: '%s'", topic, payload) + m.log.TRACE.Printf("send %s: '%s'", topic, payload) token := m.Client.Publish(topic, m.qos, false, payload) if token.WaitTimeout(publishTimeout) { return token.Error() @@ -185,7 +191,7 @@ func (m *MqttClient) BoolSetter(param, topic, message string) BoolSetter { return err } - mlog.TRACE.Printf("send %s: '%s'", topic, payload) + m.log.TRACE.Printf("send %s: '%s'", topic, payload) token := m.Client.Publish(topic, m.qos, false, payload) if token.WaitTimeout(publishTimeout) { return token.Error() @@ -199,14 +205,15 @@ func (m *MqttClient) BoolSetter(param, topic, message string) BoolSetter { func (m *MqttClient) WaitForToken(token mqtt.Token) { if token.WaitTimeout(publishTimeout) { if token.Error() != nil { - mlog.ERROR.Printf("error: %s", token.Error()) + m.log.ERROR.Printf("error: %s", token.Error()) } } else { - mlog.DEBUG.Println("timeout") + m.log.DEBUG.Println("timeout") } } type msgHandler struct { + log *api.Logger once sync.Once mux sync.Mutex updated time.Time @@ -217,7 +224,7 @@ type msgHandler struct { } func (h *msgHandler) Receive(payload string) { - mlog.TRACE.Printf("recv %s: '%s'", h.topic, payload) + h.log.TRACE.Printf("recv %s: '%s'", h.topic, payload) h.mux.Lock() defer h.mux.Unlock() @@ -231,7 +238,7 @@ func (h *msgHandler) waitForInitialValue() { defer h.mux.Unlock() if h.updated.IsZero() { - mlog.TRACE.Printf("%s wait for initial value", h.topic) + h.log.TRACE.Printf("%s wait for initial value", h.topic) // wait for initial update for h.updated.IsZero() { @@ -260,20 +267,8 @@ func (h *msgHandler) floatGetter() (float64, error) { } func (h *msgHandler) intGetter() (int64, error) { - h.once.Do(h.waitForInitialValue) - h.mux.Lock() - defer h.mux.Unlock() - - if elapsed := time.Since(h.updated); h.timeout != 0 && elapsed > h.timeout { - return 0, fmt.Errorf("%s outdated: %v", h.topic, elapsed.Truncate(time.Second)) - } - - val, err := strconv.ParseInt(h.payload, 10, 64) - if err != nil { - return 0, fmt.Errorf("%s invalid: '%s'", h.topic, h.payload) - } - - return int64(h.multiplier) * val, nil + f, err := h.floatGetter() + return int64(math.Round(f)), err } func (h *msgHandler) stringGetter() (string, error) { diff --git a/vehicle/audi.go b/vehicle/audi.go index 866f59973..1d28b871f 100644 --- a/vehicle/audi.go +++ b/vehicle/audi.go @@ -70,7 +70,7 @@ func NewAudiFromConfig(log *api.Logger, other map[string]interface{}) api.Vehicl vin: cc.VIN, } - v.chargeStateG = provider.NewCached(v.chargeState, cc.Cache).FloatGetter() + v.chargeStateG = provider.NewCached(log, v.chargeState, cc.Cache).FloatGetter() return v } diff --git a/vehicle/bmw.go b/vehicle/bmw.go index 5c27b7a32..3710bbaa4 100644 --- a/vehicle/bmw.go +++ b/vehicle/bmw.go @@ -52,7 +52,7 @@ func NewBMWFromConfig(log *api.Logger, other map[string]interface{}) api.Vehicle vin: cc.VIN, } - v.chargeStateG = provider.NewCached(v.chargeState, cc.Cache).FloatGetter() + v.chargeStateG = provider.NewCached(log, v.chargeState, cc.Cache).FloatGetter() return v } diff --git a/vehicle/nissan.go b/vehicle/nissan.go index 61ef40838..5cbafa21b 100644 --- a/vehicle/nissan.go +++ b/vehicle/nissan.go @@ -43,7 +43,7 @@ func NewNissanFromConfig(log *api.Logger, other map[string]interface{}) api.Vehi session: session, } - v.chargeStateG = provider.NewCached(v.chargeState, cc.Cache).FloatGetter() + v.chargeStateG = provider.NewCached(log, v.chargeState, cc.Cache).FloatGetter() return v } diff --git a/vehicle/tesla.go b/vehicle/tesla.go index 1f62390d2..888fb383e 100644 --- a/vehicle/tesla.go +++ b/vehicle/tesla.go @@ -57,8 +57,8 @@ func NewTeslaFromConfig(log *api.Logger, other map[string]interface{}) api.Vehic log.FATAL.Fatal("cannot create tesla: vin not found") } - v.chargeStateG = provider.NewCached(v.chargeState, cc.Cache).FloatGetter() - v.chargedEnergyG = provider.NewCached(v.chargedEnergy, cc.Cache).FloatGetter() + v.chargeStateG = provider.NewCached(log, v.chargeState, cc.Cache).FloatGetter() + v.chargedEnergyG = provider.NewCached(log, v.chargedEnergy, cc.Cache).FloatGetter() return v } diff --git a/vehicle/vehicle.go b/vehicle/vehicle.go index f3e338389..dc834e33f 100644 --- a/vehicle/vehicle.go +++ b/vehicle/vehicle.go @@ -40,7 +40,7 @@ func NewConfigurableFromConfig(log *api.Logger, other map[string]interface{}) ap getter := provider.NewFloatGetterFromConfig(log, cc.Charge) if cc.Cache > 0 { - getter = provider.NewCached(getter, cc.Cache).FloatGetter() + getter = provider.NewCached(log, getter, cc.Cache).FloatGetter() } return &Vehicle{