From 2d187750fbb61cdc2ea2df3502f433a195dcc75f Mon Sep 17 00:00:00 2001 From: andig Date: Fri, 31 Jul 2020 07:37:49 +0200 Subject: [PATCH] Don't commit errors and warnings to cache --- cmd/root.go | 8 +++--- cmd/setup.go | 5 ++-- {server => util/pipe}/limiter.go | 35 ++++++++++++++++++++++++++- {server => util/pipe}/limiter_test.go | 2 +- 4 files changed, 43 insertions(+), 7 deletions(-) rename {server => util/pipe}/limiter.go (77%) rename {server => util/pipe}/limiter_test.go (99%) diff --git a/cmd/root.go b/cmd/root.go index 10076f3c4..54df67654 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -8,14 +8,16 @@ import ( "github.com/andig/evcc/server" "github.com/andig/evcc/server/updater" "github.com/andig/evcc/util" + "github.com/andig/evcc/util/pipe" "github.com/spf13/cobra" "github.com/spf13/viper" ) var ( - log = util.NewLogger("main") - cfgFile string + ignoreParams = []string{"warn", "error", "fatal"} // don't add to cache + log = util.NewLogger("main") + cfgFile string ) // rootCmd represents the base command when called without any subcommands @@ -134,7 +136,7 @@ func run(cmd *cobra.Command, args []string) { // value cache cache := util.NewCache() - go cache.Run(tee.Attach()) + go cache.Run(pipe.NewDropper(ignoreParams...).Pipe(tee.Attach())) // setup loadpoints site := loadConfig(conf) diff --git a/cmd/setup.go b/cmd/setup.go index 429a02843..d4e39e66a 100644 --- a/cmd/setup.go +++ b/cmd/setup.go @@ -11,6 +11,7 @@ import ( "github.com/andig/evcc/push" "github.com/andig/evcc/server" "github.com/andig/evcc/util" + "github.com/andig/evcc/util/pipe" "github.com/spf13/viper" ) @@ -30,11 +31,11 @@ func configureDatabase(conf server.InfluxConfig, loadPoints []*core.LoadPoint, i ) // eliminate duplicate values - dedupe := server.NewDeduplicator(30*time.Minute, "socCharge") + dedupe := pipe.NewDeduplicator(30*time.Minute, "socCharge") in = dedupe.Pipe(in) // reduce number of values written to influx - limiter := server.NewLimiter(5 * time.Second) + limiter := pipe.NewLimiter(5 * time.Second) in = limiter.Pipe(in) go influx.Run(loadPoints, in) diff --git a/server/limiter.go b/util/pipe/limiter.go similarity index 77% rename from server/limiter.go rename to util/pipe/limiter.go index 9790bfc22..6a11f757f 100644 --- a/server/limiter.go +++ b/util/pipe/limiter.go @@ -1,4 +1,4 @@ -package server +package pipe import ( "time" @@ -100,3 +100,36 @@ func (l *Limiter) Pipe(in <-chan util.Param) <-chan util.Param { go l.pipe(in, out) return out } + +// Dropper allows filtering of channel data by given criteria +type Dropper struct { + filter []string +} + +// NewDropper creates Dropper +func NewDropper(filter ...string) Piper { + return &Dropper{filter} +} + +func (l *Dropper) pipe(in <-chan util.Param, out chan<- util.Param) { + for p := range in { + var remove bool + for _, filtered := range l.filter { + if p.Key == filtered { + remove = true + break + } + } + + if !remove { + out <- p + } + } +} + +// Pipe creates a new filtered output channel for given input channel +func (l *Dropper) 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/util/pipe/limiter_test.go similarity index 99% rename from server/limiter_test.go rename to util/pipe/limiter_test.go index 785760187..b96f801af 100644 --- a/server/limiter_test.go +++ b/util/pipe/limiter_test.go @@ -1,4 +1,4 @@ -package server +package pipe import ( "runtime"