Don't commit errors and warnings to cache
This commit is contained in:
parent
ff880489e6
commit
2d187750fb
4 changed files with 43 additions and 7 deletions
|
|
@ -1,102 +0,0 @@
|
|||
package server
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"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 util.Param) <-chan util.Param
|
||||
}
|
||||
|
||||
type cacheItem struct {
|
||||
updated time.Time
|
||||
val interface{}
|
||||
}
|
||||
|
||||
// Deduplicator allows filtering of channel data by given criteria
|
||||
type Deduplicator struct {
|
||||
clock clock.Clock
|
||||
interval time.Duration
|
||||
filter map[string]interface{}
|
||||
cache map[string]cacheItem
|
||||
}
|
||||
|
||||
// NewDeduplicator creates Deduplicator
|
||||
func NewDeduplicator(interval time.Duration, filter ...string) Piper {
|
||||
l := &Deduplicator{
|
||||
clock: clock.New(),
|
||||
interval: interval,
|
||||
filter: make(map[string]interface{}),
|
||||
cache: make(map[string]cacheItem),
|
||||
}
|
||||
|
||||
for _, f := range filter {
|
||||
l.filter[f] = struct{}{}
|
||||
}
|
||||
|
||||
return l
|
||||
}
|
||||
|
||||
func (l *Deduplicator) pipe(in <-chan util.Param, out chan<- util.Param) {
|
||||
for p := range in {
|
||||
key := p.UniqueID()
|
||||
item, cached := l.cache[key]
|
||||
_, filtered := l.filter[p.Key]
|
||||
|
||||
// forward if not cached
|
||||
if !cached || !filtered || filtered &&
|
||||
(l.clock.Since(item.updated) >= l.interval || p.Val != item.val) {
|
||||
l.cache[key] = cacheItem{updated: l.clock.Now(), val: p.Val}
|
||||
out <- p
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Pipe creates a new filtered output channel for given input channel
|
||||
func (l *Deduplicator) Pipe(in <-chan util.Param) <-chan util.Param {
|
||||
out := make(chan util.Param)
|
||||
go l.pipe(in, out)
|
||||
return out
|
||||
}
|
||||
|
||||
// Limiter allows filtering of channel data by given criteria
|
||||
type Limiter struct {
|
||||
clock clock.Clock
|
||||
interval time.Duration
|
||||
cache map[string]cacheItem
|
||||
}
|
||||
|
||||
// NewLimiter creates limiter
|
||||
func NewLimiter(interval time.Duration) Piper {
|
||||
l := &Limiter{
|
||||
clock: clock.New(),
|
||||
interval: interval,
|
||||
cache: make(map[string]cacheItem),
|
||||
}
|
||||
|
||||
return l
|
||||
}
|
||||
|
||||
func (l *Limiter) pipe(in <-chan util.Param, out chan<- util.Param) {
|
||||
for p := range in {
|
||||
key := p.UniqueID()
|
||||
item, cached := l.cache[key]
|
||||
|
||||
// forward if not cached or expired
|
||||
if !cached || l.clock.Since(item.updated) >= l.interval {
|
||||
l.cache[key] = cacheItem{updated: l.clock.Now(), val: p.Val}
|
||||
out <- p
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Pipe creates a new filtered output channel for given input channel
|
||||
func (l *Limiter) Pipe(in <-chan util.Param) <-chan util.Param {
|
||||
out := make(chan util.Param)
|
||||
go l.pipe(in, out)
|
||||
return out
|
||||
}
|
||||
|
|
@ -1,124 +0,0 @@
|
|||
package server
|
||||
|
||||
import (
|
||||
"runtime"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/andig/evcc/util"
|
||||
"github.com/benbjohnson/clock"
|
||||
)
|
||||
|
||||
func TestLimiter(t *testing.T) {
|
||||
l := NewLimiter(time.Hour).(*Limiter)
|
||||
clck := clock.NewMock()
|
||||
l.clock = clck
|
||||
|
||||
in := make(chan util.Param)
|
||||
out := l.Pipe(in)
|
||||
|
||||
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.Val = 2
|
||||
in <- p
|
||||
|
||||
runtime.Gosched()
|
||||
select {
|
||||
case o := <-out:
|
||||
t.Errorf("unexpected param %v", o)
|
||||
case <-time.After(time.Millisecond):
|
||||
}
|
||||
|
||||
clck.Add(2 * l.interval)
|
||||
p.Val = 3
|
||||
in <- p
|
||||
|
||||
if o := <-out; o.Key != p.Key || o.Val != p.Val {
|
||||
t.Errorf("unexpected param %v", o)
|
||||
}
|
||||
|
||||
// allow nils
|
||||
clck.Add(2 * l.interval)
|
||||
p.Val = nil
|
||||
in <- p
|
||||
|
||||
if o := <-out; o.Key != p.Key || o.Val != p.Val {
|
||||
t.Errorf("unexpected param %v", o)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDeduplicator(t *testing.T) {
|
||||
l := NewDeduplicator(time.Hour, "filtered").(*Deduplicator)
|
||||
clck := clock.NewMock()
|
||||
l.clock = clck
|
||||
|
||||
in := make(chan util.Param)
|
||||
out := l.Pipe(in)
|
||||
|
||||
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 = util.Param{Key: "k", Val: 2}
|
||||
in <- p
|
||||
|
||||
if o := <-out; o.Key != p.Key || o.Val != p.Val {
|
||||
t.Errorf("unexpected param %v", o)
|
||||
}
|
||||
|
||||
// allow nils
|
||||
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 = 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 = util.Param{Key: "filtered", Val: 4}
|
||||
in <- p
|
||||
|
||||
if o := <-out; o.Key != p.Key || o.Val != p.Val {
|
||||
t.Errorf("unexpected param %v", o)
|
||||
}
|
||||
|
||||
// resend
|
||||
in <- p
|
||||
|
||||
runtime.Gosched()
|
||||
select {
|
||||
case o := <-out:
|
||||
t.Errorf("unexpected param %v", o)
|
||||
case <-time.After(time.Millisecond):
|
||||
}
|
||||
|
||||
// resend later
|
||||
clck.Add(2 * l.interval)
|
||||
in <- p
|
||||
|
||||
if o := <-out; o.Key != p.Key || o.Val != p.Val {
|
||||
t.Errorf("unexpected param %v", o)
|
||||
}
|
||||
|
||||
// allow nils
|
||||
p = util.Param{Key: "filtered", Val: nil}
|
||||
in <- p
|
||||
|
||||
if o := <-out; o.Key != p.Key || o.Val != p.Val {
|
||||
t.Errorf("unexpected param %v", o)
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue