evcc-io/util/monitor.go

129 lines
2.7 KiB
Go

package util
import (
"context"
"sync"
"time"
"github.com/benbjohnson/clock"
"github.com/evcc-io/evcc/api"
)
// Monitor monitors values for regular updates
type Monitor[T any] struct {
val T
clock clock.Clock
mu sync.RWMutex
once sync.Once
done chan struct{}
updated time.Time
timeout time.Duration
}
// NewMonitor created a new monitor with given timeout
func NewMonitor[T any](timeout time.Duration) *Monitor[T] {
res := &Monitor[T]{
clock: clock.New(),
done: make(chan struct{}),
timeout: timeout,
}
return res
}
// WithClock sets the a clock for debugging
func (m *Monitor[T]) WithClock(clock clock.Clock) *Monitor[T] {
m.clock = clock
return m
}
// Set updates the current value and timestamp
func (m *Monitor[T]) Set(val T) {
m.SetFunc(func(_ T) T { return val })
}
// SetFunc updates the current value and timestamp while holding the lock
func (m *Monitor[T]) SetFunc(set func(T) T) {
m.mu.Lock()
defer m.mu.Unlock()
m.val = set(m.val)
m.updated = m.clock.Now()
m.once.Do(func() { close(m.done) })
}
// Get returns the current value or ErrOutdated if timeout exceeded
func (m *Monitor[T]) Get() (T, error) {
return m.GetContext(context.Background())
}
// GetContext returns the current value or ErrOutdated if timeout exceeded.
// The context can cancel the blocking first-call wait.
func (m *Monitor[T]) GetContext(ctx context.Context) (T, error) {
var res T
err := m.GetFuncContext(ctx, func(v T) {
res = v
})
return res, err
}
// GetFunc returns the current value or ErrOutdated if timeout exceeded while holding the lock
func (m *Monitor[T]) GetFunc(get func(T)) error {
return m.GetFuncContext(context.Background(), get)
}
// GetFuncContext returns the current value or ErrOutdated if timeout exceeded while holding the lock.
// The context can cancel the blocking first-call wait.
func (m *Monitor[T]) GetFuncContext(ctx context.Context, get func(T)) error {
m.mu.RLock()
defer m.mu.RUnlock()
// without timeout set, error if not yet received
if m.timeout == 0 {
select {
case <-m.done:
get(m.val)
return nil
default:
return api.ErrOutdated
}
}
if m.clock.Since(m.updated) > m.timeout {
err := api.ErrOutdated
// wait once on very first call
if m.updated.IsZero() {
m.mu.RUnlock()
// mark as waited once
m.mu.Lock()
// TODO fix and test
m.updated = m.updated.Add(time.Nanosecond)
m.mu.Unlock()
select {
case <-m.done:
// got value and updated timestamp
err = nil
case <-ctx.Done():
err = ctx.Err()
case <-m.clock.After(m.timeout):
}
m.mu.RLock()
}
get(m.val)
return err
}
get(m.val)
return nil
}
// Done signals if monitor has been updated at least once
func (m *Monitor[T]) Done() <-chan struct{} {
return m.done
}