Plugins: make watchdog deferable (#26790)
This commit is contained in:
parent
466a23d474
commit
f3c87f00b4
2 changed files with 217 additions and 21 deletions
|
|
@ -8,18 +8,21 @@ import (
|
|||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/benbjohnson/clock"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
)
|
||||
|
||||
type watchdogPlugin struct {
|
||||
mu sync.Mutex
|
||||
ctx context.Context
|
||||
log *util.Logger
|
||||
reset []string
|
||||
initial *string
|
||||
set Config
|
||||
timeout time.Duration
|
||||
cancel func()
|
||||
mu sync.Mutex
|
||||
ctx context.Context
|
||||
log *util.Logger
|
||||
reset []string
|
||||
initial *string
|
||||
set Config
|
||||
timeout time.Duration
|
||||
deferred bool
|
||||
cancel func()
|
||||
clock clock.Clock
|
||||
}
|
||||
|
||||
func init() {
|
||||
|
|
@ -33,6 +36,7 @@ func NewWatchDogFromConfig(ctx context.Context, other map[string]any) (Plugin, e
|
|||
Initial *string
|
||||
Set Config
|
||||
Timeout time.Duration
|
||||
Defer bool `mapstructure:"defer"`
|
||||
}
|
||||
|
||||
if err := util.DecodeOther(other, &cc); err != nil {
|
||||
|
|
@ -40,12 +44,14 @@ func NewWatchDogFromConfig(ctx context.Context, other map[string]any) (Plugin, e
|
|||
}
|
||||
|
||||
o := &watchdogPlugin{
|
||||
ctx: ctx,
|
||||
log: util.ContextLoggerWithDefault(ctx, util.NewLogger("watchdog")),
|
||||
reset: cc.Reset,
|
||||
initial: cc.Initial,
|
||||
set: cc.Set,
|
||||
timeout: cc.Timeout,
|
||||
ctx: ctx,
|
||||
log: util.ContextLoggerWithDefault(ctx, util.NewLogger("watchdog")),
|
||||
reset: cc.Reset,
|
||||
initial: cc.Initial,
|
||||
set: cc.Set,
|
||||
timeout: cc.Timeout,
|
||||
deferred: cc.Defer,
|
||||
clock: clock.New(),
|
||||
}
|
||||
|
||||
return o, nil
|
||||
|
|
@ -64,20 +70,35 @@ func (o *watchdogPlugin) wdt(ctx context.Context, set func() error) {
|
|||
}
|
||||
}
|
||||
|
||||
type deferredState[T comparable] struct {
|
||||
val T
|
||||
timer *clock.Timer
|
||||
}
|
||||
|
||||
// setter is the generic setter function for watchdogPlugin
|
||||
// it is currently not possible to write this as a method
|
||||
func setter[T comparable](o *watchdogPlugin, set func(T) error, reset []T) func(T) error {
|
||||
return func(val T) error {
|
||||
o.mu.Lock()
|
||||
defer o.mu.Unlock()
|
||||
var state *deferredState[T]
|
||||
var lastUpdated time.Time
|
||||
var last *T
|
||||
|
||||
// stop wdt on new write
|
||||
// stop running wdt
|
||||
stopWdt := func() {
|
||||
if o.cancel != nil {
|
||||
o.cancel()
|
||||
o.cancel = nil
|
||||
}
|
||||
}
|
||||
|
||||
// start wdt on non-reset value
|
||||
// set value and start wdt
|
||||
setAndStartWdt := func(val T) error {
|
||||
if err := set(val); err != nil {
|
||||
return err
|
||||
}
|
||||
lastUpdated = o.clock.Now()
|
||||
last = &val
|
||||
|
||||
// start wdt for non-reset value
|
||||
if !slices.Contains(reset, val) {
|
||||
var ctx context.Context
|
||||
ctx, o.cancel = context.WithCancel(context.Background())
|
||||
|
|
@ -86,11 +107,65 @@ func setter[T comparable](o *watchdogPlugin, set func(T) error, reset []T) func(
|
|||
o.mu.Lock()
|
||||
defer o.mu.Unlock()
|
||||
|
||||
return set(val)
|
||||
if err := set(val); err != nil {
|
||||
return err
|
||||
}
|
||||
lastUpdated = o.clock.Now()
|
||||
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
return set(val)
|
||||
return nil
|
||||
}
|
||||
|
||||
return func(val T) error {
|
||||
o.mu.Lock()
|
||||
defer o.mu.Unlock()
|
||||
|
||||
// cancel deferred update
|
||||
if state != nil {
|
||||
state.timer.Stop()
|
||||
state = nil
|
||||
}
|
||||
|
||||
// if value unchanged, let wdt continue running
|
||||
// TODO refactor use of last value once batterymode is set only once, currently required to avoid defer loops
|
||||
if last != nil && *last == val && o.cancel != nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
// calculate remaining deferred delay
|
||||
delay := max(0, o.timeout+5*time.Second-o.clock.Since(lastUpdated))
|
||||
|
||||
// defer update to non-reset value
|
||||
if o.deferred && delay > 0 && !lastUpdated.IsZero() && !slices.Contains(reset, val) {
|
||||
stopWdt()
|
||||
|
||||
// store deferred value
|
||||
state = &deferredState[T]{
|
||||
val: val,
|
||||
timer: o.clock.AfterFunc(delay, func() {
|
||||
o.mu.Lock()
|
||||
defer o.mu.Unlock()
|
||||
|
||||
state = nil
|
||||
|
||||
o.log.TRACE.Printf("deferred update executing: to=%v", val)
|
||||
if err := setAndStartWdt(val); err != nil {
|
||||
o.log.ERROR.Printf("deferred update failed: %v", err)
|
||||
return
|
||||
}
|
||||
o.log.TRACE.Printf("deferred update completed: value=%v", val)
|
||||
}),
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// immediate update
|
||||
stopWdt()
|
||||
return setAndStartWdt(val)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ import (
|
|||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/benbjohnson/clock"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/stretchr/testify/require"
|
||||
"golang.org/x/sync/errgroup"
|
||||
|
|
@ -16,6 +17,7 @@ func TestWatchdogSetterConcurrency(t *testing.T) {
|
|||
p := &watchdogPlugin{
|
||||
log: util.NewLogger("foo"),
|
||||
timeout: 10 * time.Nanosecond,
|
||||
clock: clock.New(),
|
||||
}
|
||||
|
||||
var u atomic.Uint32
|
||||
|
|
@ -44,3 +46,122 @@ func TestWatchdogSetterConcurrency(t *testing.T) {
|
|||
|
||||
require.NoError(t, eg.Wait())
|
||||
}
|
||||
|
||||
func TestWatchdogDeferredUpdate(t *testing.T) {
|
||||
// Test: Value 1 → 3 → 2 with delay
|
||||
// 1 → 3: delayed (target 3 is non-reset)
|
||||
// 3 → 2: delayed (target 2 is non-reset)
|
||||
// Expected: [1, <delay>, 3, <delay>, 2]
|
||||
|
||||
timeout := 60 * time.Second
|
||||
c := clock.NewMock()
|
||||
p := &watchdogPlugin{
|
||||
log: util.NewLogger("test"),
|
||||
timeout: timeout,
|
||||
deferred: true,
|
||||
clock: c,
|
||||
}
|
||||
|
||||
var calls []int
|
||||
set := setter(p, func(i int) error {
|
||||
calls = append(calls, i)
|
||||
return nil
|
||||
}, []int{1}) // 1 is reset value
|
||||
|
||||
// Value 1 (reset) → should set immediately
|
||||
require.NoError(t, set(1))
|
||||
require.Equal(t, []int{1}, calls)
|
||||
|
||||
// Value 3 (target is non-reset) → should be delayed
|
||||
require.NoError(t, set(3))
|
||||
require.Equal(t, []int{1}, calls, "Value 3 should not be set yet")
|
||||
|
||||
// Wait for delay
|
||||
expectedDelay := p.timeout + 5*time.Second
|
||||
c.Add(expectedDelay)
|
||||
|
||||
// Now value 3 should be set
|
||||
require.Equal(t, []int{1, 3}, calls)
|
||||
|
||||
// Value 2 (non-reset to non-reset) → should delay
|
||||
require.NoError(t, set(2))
|
||||
require.Equal(t, []int{1, 3}, calls, "Value 2 should not be set yet")
|
||||
|
||||
// Wait for delay
|
||||
c.Add(expectedDelay)
|
||||
|
||||
// Now value 2 should be set (exactly once)
|
||||
require.Equal(t, []int{1, 3, 2}, calls)
|
||||
}
|
||||
|
||||
func TestWatchdogCancelPendingDeferredUpdate(t *testing.T) {
|
||||
// Test: Value 3 → 2 started, then set Value 1 during delay
|
||||
// Expected: Deferred update cancelled, Value 1 set immediately
|
||||
|
||||
timeout := 60 * time.Second
|
||||
c := clock.NewMock()
|
||||
p := &watchdogPlugin{
|
||||
log: util.NewLogger("test"),
|
||||
timeout: timeout,
|
||||
deferred: true,
|
||||
clock: c,
|
||||
}
|
||||
|
||||
var calls []int
|
||||
set := setter(p, func(i int) error {
|
||||
calls = append(calls, i)
|
||||
return nil
|
||||
}, []int{1}) // 1 is reset value
|
||||
|
||||
// Value 3 (non-reset)
|
||||
require.NoError(t, set(3))
|
||||
require.Equal(t, []int{3}, calls)
|
||||
|
||||
// Value 2 (deferred update)
|
||||
require.NoError(t, set(2))
|
||||
require.Equal(t, []int{3}, calls, "Value 2 should not be set yet")
|
||||
|
||||
// Wait a bit but not the full delay
|
||||
c.Add(30 * time.Second)
|
||||
|
||||
// Value 1 (reset) → should cancel pending deferred update and set immediately
|
||||
require.NoError(t, set(1))
|
||||
require.Equal(t, []int{3, 1}, calls, "Value 1 should be set, Value 2 should be cancelled")
|
||||
|
||||
// Wait for what would have been the original delay
|
||||
c.Add(timeout + 5*time.Second)
|
||||
|
||||
// Value 2 should still not have been set
|
||||
require.Equal(t, []int{3, 1}, calls, "Value 2 should remain cancelled")
|
||||
}
|
||||
|
||||
func TestWatchdogDelayBackwardCompatibility(t *testing.T) {
|
||||
// Test: deferred=false behaves like old implementation
|
||||
// Expected: All updates immediate
|
||||
|
||||
p := &watchdogPlugin{
|
||||
log: util.NewLogger("test"),
|
||||
timeout: 60 * time.Second,
|
||||
deferred: false, // explicitly false
|
||||
clock: clock.New(),
|
||||
}
|
||||
|
||||
var calls []int
|
||||
set := setter(p, func(i int) error {
|
||||
calls = append(calls, i)
|
||||
return nil
|
||||
}, []int{1}) // 1 is reset value
|
||||
|
||||
// All updates should be immediate
|
||||
require.NoError(t, set(1))
|
||||
require.Equal(t, []int{1}, calls)
|
||||
|
||||
require.NoError(t, set(3))
|
||||
require.Equal(t, []int{1, 3}, calls)
|
||||
|
||||
require.NoError(t, set(2))
|
||||
require.Equal(t, []int{1, 3, 2}, calls, "Value 2 should be set immediately (no delay)")
|
||||
|
||||
require.NoError(t, set(4))
|
||||
require.Equal(t, []int{1, 3, 2, 4}, calls, "Value 4 should be set immediately (no delay)")
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue