diff --git a/plugin/watchdog.go b/plugin/watchdog.go index 8c4a579c8..a192d819d 100644 --- a/plugin/watchdog.go +++ b/plugin/watchdog.go @@ -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) } } diff --git a/plugin/watchdog_test.go b/plugin/watchdog_test.go index 52d0274e9..be9dc4151 100644 --- a/plugin/watchdog_test.go +++ b/plugin/watchdog_test.go @@ -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, , 3, , 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)") +}