working on valuemonitor
This commit is contained in:
+70
-8
@@ -1,8 +1,10 @@
|
||||
package settings
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -24,11 +26,67 @@ var (
|
||||
)
|
||||
|
||||
type Settings struct {
|
||||
mutex sync.RWMutex
|
||||
contents map[string]string
|
||||
filename string
|
||||
modtime time.Time
|
||||
LogUpdates atomic.Bool
|
||||
mutex sync.RWMutex
|
||||
context context.Context
|
||||
contents map[string]string
|
||||
monitorKeys map[string]chan<- string
|
||||
filename string
|
||||
modtime time.Time
|
||||
LogUpdates atomic.Bool
|
||||
}
|
||||
|
||||
func insertOrTimeout(ctx context.Context, c chan<- string, s string) error {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case c <- s:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Settings) RegisterKeyToMonitor(key string, updateOutputChannel chan<- string) error {
|
||||
s.mutex.Lock()
|
||||
defer s.mutex.Unlock()
|
||||
if s.monitorKeys == nil {
|
||||
s.monitorKeys = make(map[string]chan<- string, 1)
|
||||
}
|
||||
if _, found := s.monitorKeys[key]; found {
|
||||
return fmt.Errorf("key %q already in monitorKeys map", key)
|
||||
}
|
||||
if err := insertOrTimeout(s.context, updateOutputChannel, s.contents[key]); err != nil {
|
||||
return fmt.Errorf("error inserting: %w", err)
|
||||
}
|
||||
s.monitorKeys[key] = updateOutputChannel
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Settings) Dump() ([]byte, error) {
|
||||
var buffer bytes.Buffer
|
||||
if err := s.DumpToWriter(&buffer); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
func (s *Settings) DumpToWriter(w io.Writer) error {
|
||||
writer := bufio.NewWriter(w)
|
||||
s.mutex.RLock()
|
||||
defer s.mutex.RUnlock()
|
||||
var index int
|
||||
for k, v := range s.contents {
|
||||
writer.WriteString(fmt.Sprintf("%s=%s", k, v))
|
||||
if index < len(s.contents)-1 {
|
||||
writer.WriteByte('\n')
|
||||
}
|
||||
index++
|
||||
}
|
||||
return writer.Flush()
|
||||
}
|
||||
|
||||
func (s *Settings) DumpToJson() ([]byte, error) {
|
||||
s.mutex.RLock()
|
||||
defer s.mutex.RUnlock()
|
||||
return json.Marshal(s.contents)
|
||||
}
|
||||
|
||||
func (s *Settings) GetFilename() string {
|
||||
@@ -68,6 +126,9 @@ func (s *Settings) lockedSetKeyValue(key, value string) bool {
|
||||
log.Printf("Settings::lockedKeyValue %s: storing key '%s' with value '%s'\n", baseFilename, key, value)
|
||||
}
|
||||
s.contents[strings.Clone(key)] = strings.Clone(value)
|
||||
if updateChannel, found := s.monitorKeys[key]; found {
|
||||
insertOrTimeout(s.context, updateChannel, value)
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -132,10 +193,10 @@ func (s *Settings) Update() error {
|
||||
return s.lockedUpdate()
|
||||
}
|
||||
|
||||
func (s *Settings) maintenanceRoutine(ctx context.Context) {
|
||||
func (s *Settings) maintenanceRoutine() {
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
case <-s.context.Done():
|
||||
return
|
||||
case <-time.After(MaintenanceRoutinePace):
|
||||
err := s.Update()
|
||||
@@ -154,6 +215,7 @@ func NewSettings(filename string, logUpdates bool) *Settings {
|
||||
|
||||
func NewSettingsWithContext(ctx context.Context, filename string, logUpdates bool) *Settings {
|
||||
settings := &Settings{
|
||||
context: ctx,
|
||||
contents: make(map[string]string),
|
||||
filename: strings.Clone(filename),
|
||||
}
|
||||
@@ -161,7 +223,7 @@ func NewSettingsWithContext(ctx context.Context, filename string, logUpdates boo
|
||||
if err := settings.lockedUpdate(); err != nil && (!errors.Is(err, os.ErrNotExist) || WarnIfFileNotFound) {
|
||||
log.Printf("NewSettingsWithContext %s error from initial update: %v\n", filepath.Base(filename), err)
|
||||
}
|
||||
go settings.maintenanceRoutine(ctx)
|
||||
go settings.maintenanceRoutine()
|
||||
return settings
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user