129 lines
2.7 KiB
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
|
|
}
|