diff --git a/cmd/root.go b/cmd/root.go index 89ee781d3..4476c4c86 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -128,25 +128,13 @@ func checkVersion() { } } -var teeIsChained bool // controles piping of first channel in teed chain - -// tee splits a tee channel from an input channel and starts a goroutine -// for duplicateing channel values to replacement input and tee channel -func tee(in chan core.Param) (chan core.Param, <-chan core.Param) { - gen := make(chan core.Param) - tee := make(chan core.Param) - - go func(teeIsChained bool) { - for i := range gen { - if teeIsChained { - in <- i - } - tee <- i +// handle UI update requests +func handleUI(triggerChan <-chan struct{}, loadPoints []*core.LoadPoint) { + for range triggerChan { + for _, lp := range loadPoints { + lp.Update() } - }(teeIsChained) - - teeIsChained = true - return gen, tee + } } func run(cmd *cobra.Command, args []string) { @@ -174,8 +162,7 @@ func run(cmd *cobra.Command, args []string) { loadPoints := loadConfig(conf, notificationChan) // start broadcasting values - valueChan := make(chan core.Param) - triggerChan := make(chan struct{}) + tee := &Tee{} // setup influx if viper.Get("influx") != nil { @@ -187,28 +174,32 @@ func run(cmd *cobra.Command, args []string) { conf.Influx.Password, ) - var teeChan <-chan core.Param - valueChan, teeChan = tee(valueChan) - // eliminate duplicate values dedupe := server.NewDeduplicator(30*time.Minute, "socCharge") - teeChan = dedupe.Pipe(teeChan) + pipeChan := dedupe.Pipe(tee.Attach()) // reduce number of values written to influx limiter := server.NewLimiter(5 * time.Second) - teeChan = limiter.Pipe(teeChan) + pipeChan = limiter.Pipe(pipeChan) - go influx.Run(teeChan) + go influx.Run(pipeChan) } // create webserver socketHub := server.NewSocketHub() httpd := server.NewHTTPd(uri, conf.Menu, loadPoints[0], socketHub) - var teeChan <-chan core.Param - valueChan, teeChan = tee(valueChan) + triggerChan := make(chan struct{}) - go socketHub.Run(teeChan, triggerChan) + // handle UI update requests whenever browser connects + go handleUI(triggerChan, loadPoints) + + // publish to UI + go socketHub.Run(tee.Attach(), triggerChan) + + // setup values channel + valueChan := make(chan util.Param) + go tee.Run(valueChan) // start all loadpoints for _, lp := range loadPoints { @@ -217,14 +208,5 @@ func run(cmd *cobra.Command, args []string) { go lp.Run(conf.Interval) } - // handle UI update requests whenever browser connects - go func() { - for range triggerChan { - for _, lp := range loadPoints { - lp.Update() - } - } - }() - log.FATAL.Println(httpd.ListenAndServe()) } diff --git a/cmd/tee.go b/cmd/tee.go new file mode 100644 index 000000000..f1daaf656 --- /dev/null +++ b/cmd/tee.go @@ -0,0 +1,29 @@ +package cmd + +import "github.com/andig/evcc/util" + +// Tee distributed parameters to subscribers +type Tee struct { + recv []chan<- util.Param +} + +// Attach creates a new receiver channel and attaches it to the tee +func (t *Tee) Attach() <-chan util.Param { + out := make(chan util.Param) + t.Add(out) + return out +} + +// Add attaches a receiver channel to the tee +func (t *Tee) Add(out chan<- util.Param) { + t.recv = append(t.recv, out) +} + +// Run starts parameter distribution +func (t *Tee) Run(in <-chan util.Param) { + for msg := range in { + for _, recv := range t.recv { + recv <- msg + } + } +} diff --git a/core/api.go b/core/api.go index e1a28c919..64361c3bb 100644 --- a/core/api.go +++ b/core/api.go @@ -14,13 +14,6 @@ func (n *nilVal) String() string { return "—" } -// Param is the broadcast channel data type -type Param struct { - LoadPoint string - Key string - Val interface{} -} - // Configuration is the loadpoint feature structure type Configuration struct { Mode string `json:"mode"` diff --git a/core/loadpoint.go b/core/loadpoint.go index d32b7a06b..bdbb817e7 100644 --- a/core/loadpoint.go +++ b/core/loadpoint.go @@ -50,7 +50,7 @@ type Config struct { Enable, Disable ThresholdConfig } -// ThresholdConfig defines enable/disable hysteresis paramters +// ThresholdConfig defines enable/disable hysteresis parameters type ThresholdConfig struct { Delay time.Duration Threshold float64 @@ -72,7 +72,7 @@ type LoadPoint struct { bus evbus.Bus // event bus triggerChan chan struct{} // API updates notificationChan chan<- push.Event // notifications - uiChan chan<- Param // client push messages + uiChan chan<- util.Param // client push messages Config `mapstructure:",squash"` // exposed public configuration ChargerHandler `mapstructure:",squash"` // handle charger state and current @@ -169,7 +169,7 @@ func (lp *LoadPoint) notify(event string, attributes map[string]interface{}) { // publish sends values to UI and databases func (lp *LoadPoint) publish(key string, val interface{}) { - lp.uiChan <- Param{ + lp.uiChan <- util.Param{ LoadPoint: lp.Name, Key: key, Val: val, @@ -228,7 +228,7 @@ func (lp *LoadPoint) evChargeCurrentHandler(current int64) { } // Prepare loadpoint configuration by adding missing helper elements -func (lp *LoadPoint) Prepare(uiChan chan<- Param, notificationChan chan<- push.Event) { +func (lp *LoadPoint) Prepare(uiChan chan<- util.Param, notificationChan chan<- push.Event) { lp.notificationChan = notificationChan lp.uiChan = uiChan diff --git a/core/loadpoint_test.go b/core/loadpoint_test.go index 2bd8975ba..3caec8600 100644 --- a/core/loadpoint_test.go +++ b/core/loadpoint_test.go @@ -10,6 +10,7 @@ import ( "github.com/andig/evcc/mock" "github.com/andig/evcc/provider" "github.com/andig/evcc/push" + "github.com/andig/evcc/util" "github.com/benbjohnson/clock" "github.com/golang/mock/gomock" ) @@ -66,7 +67,7 @@ func newLoadPoint(charger api.Charger, pv, gm, cm api.Meter) *LoadPoint { lp.chargeMeter = cm } - uiChan := make(chan Param) + uiChan := make(chan util.Param) notificationChan := make(chan push.Event) lp.Prepare(uiChan, notificationChan) diff --git a/server/influxdb.go b/server/influxdb.go index 36473016b..9602769d1 100644 --- a/server/influxdb.go +++ b/server/influxdb.go @@ -4,7 +4,6 @@ import ( "sync" "time" - "github.com/andig/evcc/core" "github.com/andig/evcc/util" influxdb "github.com/influxdata/influxdb1-client/v2" ) @@ -131,7 +130,7 @@ func (m *Influx) asyncWriter(exit <-chan struct{}) <-chan struct{} { } // Run Influx publisher -func (m *Influx) Run(in <-chan core.Param) { +func (m *Influx) Run(in <-chan util.Param) { exit := make(chan struct{}) // exit signals to stop writer done := m.asyncWriter(exit) // done signals writer stopped diff --git a/server/limiter.go b/server/limiter.go index cd011636f..bc00f298a 100644 --- a/server/limiter.go +++ b/server/limiter.go @@ -3,13 +3,13 @@ package server import ( "time" - "github.com/andig/evcc/core" + "github.com/andig/evcc/util" "github.com/benbjohnson/clock" ) // Piper is the interface that data flow plugins must implement type Piper interface { - Pipe(in <-chan core.Param) <-chan core.Param + Pipe(in <-chan util.Param) <-chan util.Param } type cacheItem struct { @@ -41,7 +41,7 @@ func NewDeduplicator(interval time.Duration, filter ...string) Piper { return l } -func (l *Deduplicator) pipe(in <-chan core.Param, out chan<- core.Param) { +func (l *Deduplicator) pipe(in <-chan util.Param, out chan<- util.Param) { for p := range in { // use loadpoint + param.Key as lookup key to value cache key := p.LoadPoint + "." + p.Key @@ -58,8 +58,8 @@ func (l *Deduplicator) pipe(in <-chan core.Param, out chan<- core.Param) { } // Pipe creates a new filtered output channel for given input channel -func (l *Deduplicator) Pipe(in <-chan core.Param) <-chan core.Param { - out := make(chan core.Param) +func (l *Deduplicator) Pipe(in <-chan util.Param) <-chan util.Param { + out := make(chan util.Param) go l.pipe(in, out) return out } @@ -82,7 +82,7 @@ func NewLimiter(interval time.Duration) Piper { return l } -func (l *Limiter) pipe(in <-chan core.Param, out chan<- core.Param) { +func (l *Limiter) pipe(in <-chan util.Param, out chan<- util.Param) { for p := range in { // use loadpoint + param.Key as lookup key to value cache key := p.LoadPoint + "." + p.Key @@ -97,8 +97,8 @@ func (l *Limiter) pipe(in <-chan core.Param, out chan<- core.Param) { } // Pipe creates a new filtered output channel for given input channel -func (l *Limiter) Pipe(in <-chan core.Param) <-chan core.Param { - out := make(chan core.Param) +func (l *Limiter) Pipe(in <-chan util.Param) <-chan util.Param { + out := make(chan util.Param) go l.pipe(in, out) return out } diff --git a/server/limiter_test.go b/server/limiter_test.go index 98737074b..785760187 100644 --- a/server/limiter_test.go +++ b/server/limiter_test.go @@ -5,7 +5,7 @@ import ( "testing" "time" - "github.com/andig/evcc/core" + "github.com/andig/evcc/util" "github.com/benbjohnson/clock" ) @@ -14,10 +14,10 @@ func TestLimiter(t *testing.T) { clck := clock.NewMock() l.clock = clck - in := make(chan core.Param) + in := make(chan util.Param) out := l.Pipe(in) - p := core.Param{Key: "k", Val: 1} + p := util.Param{Key: "k", Val: 1} in <- p if o := <-out; o.Key != p.Key || o.Val != p.Val { @@ -57,17 +57,17 @@ func TestDeduplicator(t *testing.T) { clck := clock.NewMock() l.clock = clck - in := make(chan core.Param) + in := make(chan util.Param) out := l.Pipe(in) - p := core.Param{Key: "k", Val: 1} + p := util.Param{Key: "k", Val: 1} in <- p if o := <-out; o.Key != p.Key || o.Val != p.Val { t.Errorf("unexpected param %v", o) } - p = core.Param{Key: "k", Val: 2} + p = util.Param{Key: "k", Val: 2} in <- p if o := <-out; o.Key != p.Key || o.Val != p.Val { @@ -75,21 +75,21 @@ func TestDeduplicator(t *testing.T) { } // allow nils - p = core.Param{Key: "k", Val: nil} + p = util.Param{Key: "k", Val: nil} in <- p if o := <-out; o.Key != p.Key || o.Val != p.Val { t.Errorf("unexpected param %v", o) } - p = core.Param{Key: "filtered", Val: 3} + p = util.Param{Key: "filtered", Val: 3} in <- p if o := <-out; o.Key != p.Key || o.Val != p.Val { t.Errorf("unexpected param %v", o) } - p = core.Param{Key: "filtered", Val: 4} + p = util.Param{Key: "filtered", Val: 4} in <- p if o := <-out; o.Key != p.Key || o.Val != p.Val { @@ -115,7 +115,7 @@ func TestDeduplicator(t *testing.T) { } // allow nils - p = core.Param{Key: "filtered", Val: nil} + p = util.Param{Key: "filtered", Val: nil} in <- p if o := <-out; o.Key != p.Key || o.Val != p.Val { diff --git a/server/socket.go b/server/socket.go index 120bc36a2..d9d2100a4 100644 --- a/server/socket.go +++ b/server/socket.go @@ -5,7 +5,7 @@ import ( "net/http" "time" - "github.com/andig/evcc/core" + "github.com/andig/evcc/util" "github.com/gorilla/websocket" ) @@ -98,7 +98,7 @@ func encode(v interface{}) (string, error) { return s, nil } -func (h *SocketHub) broadcast(i core.Param) { +func (h *SocketHub) broadcast(i util.Param) { if len(h.clients) > 0 { val, err := encode(i.Val) if err != nil { @@ -118,7 +118,7 @@ func (h *SocketHub) broadcast(i core.Param) { } // Run starts data and status distribution -func (h *SocketHub) Run(in <-chan core.Param, triggerChan chan<- struct{}) { +func (h *SocketHub) Run(in <-chan util.Param, triggerChan chan<- struct{}) { for { select { case client := <-h.register: diff --git a/util/param.go b/util/param.go new file mode 100644 index 000000000..b5ff544f7 --- /dev/null +++ b/util/param.go @@ -0,0 +1,8 @@ +package util + +// Param is the broadcast channel data type +type Param struct { + LoadPoint string + Key string + Val interface{} +}