Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 27 additions & 4 deletions internal/filedata/document.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,14 +25,37 @@ type Document struct {
Segments *map[string]ldmodel.Segment
}

// ReadError indicates that one of the source files could not be read or parsed. It
// distinguishes a per-file failure from a failure to merge the files' contents.
type ReadError struct {
Err error
Path string
}

func (e *ReadError) Error() string {
return fmt.Sprintf("%s [%s]", e.Err, e.Path)
}

func (e *ReadError) Unwrap() error {
return e.Err
}

// ReadFile reads and parses a single data file, which may be in JSON or YAML format.
func ReadFile(path string) (Document, error) {
rawData, err := os.ReadFile(path) //nolint:gosec // G304: ok to read file into variable
if err != nil {
return Document{}, wrapReadError(err)
}
return parseDocument(rawData)
}

func wrapReadError(err error) error {
return fmt.Errorf("unable to read file: %s", err)
}

func parseDocument(rawData []byte) (Document, error) {
var data Document
var rawData []byte
var err error
if rawData, err = os.ReadFile(path); err != nil { //nolint:gosec // G304: ok to read file into variable
return data, fmt.Errorf("unable to read file: %s", err)
}
if detectJSON(rawData) {
err = json.Unmarshal(rawData, &data)
} else {
Expand Down
288 changes: 288 additions & 0 deletions internal/filedata/reloader.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,288 @@
package filedata

import (
"bytes"
"crypto/sha256"
"os"
"sync"
"time"

"github.com/launchdarkly/go-sdk-common/v3/ldlog"
)

const (
// DefaultDebounceDelay is a settle window long enough to coalesce the burst of change
// notifications produced by a single file edit, and short enough to stay responsive.
DefaultDebounceDelay = 100 * time.Millisecond
// DefaultRetryDelay bounds how long a failed reload can go uncorrected when no further
// change notification arrives, for example when the failure came from reading a file
// mid-write. Reading a local file is cheap, so this can be short.
DefaultRetryDelay = time.Second
)

// ReloaderConfig configures a Reloader.
type ReloaderConfig struct {
// Paths is the list of files to load, already resolved to absolute paths. The order is
// significant: it determines which file wins under the duplicate-key handling.
Paths []string
// DuplicateKeysHandling determines what happens when the same key appears in more than
// one file.
DuplicateKeysHandling DuplicateKeysHandling
// Loggers receives log output about reloads and failures.
Loggers ldlog.Loggers
// Apply is invoked with each successfully merged result. Calls are serialized on the
// reloader's own goroutine (or, for ReloadNow, under the same lock), so implementations
// do not need their own synchronization against other reloads. Apply and OnError must
// not call back into Close.
Apply func(MergeResult)
// OnError is invoked when a reload fails, once per distinct failure: with automatic
// retries, repeats of an identical failure do not re-invoke it (a success re-arms it).
// The error is a *ReadError when a file could not be read or parsed, or a merge error
// otherwise. The reloader logs failures itself, so implementations only need to update
// their own state.
OnError func(err error)
// DebounceDelay is how long to wait after a Trigger call for further calls to settle
// before reloading, coalescing bursts of change notifications into one reload. If zero
// or negative, each Trigger reloads immediately.
DebounceDelay time.Duration
// RetryDelay is how long to wait after a failed reload before automatically retrying,
// so that a failure observed while a file was being rewritten recovers even if no
// further change notification arrives. If zero or negative, there is no automatic
// retry.
RetryDelay time.Duration
// SkipUnchanged, if true, suppresses the Apply call when the files' raw contents are
// byte-identical to the last successfully applied contents.
SkipUnchanged bool
}

// Reloader owns the reload cycle for a set of data files: it serializes reloads, debounces
// change signals, retains the last good result on failure (by not calling Apply), retries
// after failures, and skips no-op applications.
type Reloader struct {
cfg ReloaderConfig
triggerCh chan struct{}
// retryRequestCh asks the run goroutine to arm the retry timer without performing an
// immediate reload; it is used when a synchronous ReloadNow fails.
retryRequestCh chan struct{}
closeCh chan struct{}
closeOnce sync.Once
startOnce sync.Once

// reloadMu serializes the actual load work between ReloadNow and the run goroutine.
reloadMu sync.Mutex
lastGoodHash []byte
lastErrorMsg string
}

// NewReloader creates a Reloader. The caller should perform the initial load with
// ReloadNow, route change signals to Trigger, and call Close when finished.
//
// The worker goroutine starts on the first ReloadNow or Trigger call rather than here:
// a Reloader can be constructed by a component whose lifecycle never uses it (a file data
// source built as a one-shot initializer, whose interface has no Close), and construction
// alone must not leak a goroutine.
func NewReloader(cfg ReloaderConfig) *Reloader {
return &Reloader{
cfg: cfg,
triggerCh: make(chan struct{}, 1),
retryRequestCh: make(chan struct{}, 1),
closeCh: make(chan struct{}),
}
}

func (r *Reloader) ensureStarted() {
if r.isClosed() {
return
}
r.startOnce.Do(func() {
go r.run()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are you guaranteed r.run() has executed by the time ensureStarted() returns? I recall there being some go runtime versions that did jump to execute the target of the go func and some that did not. Not sure what the current state of that is. Perhaps it doesn't matter either way due to the arrangement.

})
}

// ReloadNow synchronously loads the files and applies the result (or reports the failure).
// It is used for the initial load; a failure here schedules the same automatic retry as a
// failed triggered reload.
func (r *Reloader) ReloadNow() {
r.ensureStarted()
if !r.reload() && r.cfg.RetryDelay > 0 {
select {
case r.retryRequestCh <- struct{}{}:
default:
}
}
}

// Trigger signals that the files may have changed and a reload should happen after the
// debounce delay. It never blocks; signals arriving while a reload is already pending are
// coalesced.
func (r *Reloader) Trigger() {
r.ensureStarted()
select {
case r.triggerCh <- struct{}{}:
default:
}
}

// Close stops the reloader. It does not wait for a reload that is already in progress —
// one wedged in a blocking file read (an unresponsive network mount, for example) must not
// be able to wedge shutdown — so such a reload may still deliver its result through Apply
// or OnError shortly after Close returns; consumers must tolerate that, as they always
// have for late reloads. A reload that has not yet reached its callbacks when Close is
// called will not invoke them.
func (r *Reloader) Close() {
r.closeOnce.Do(func() {
close(r.closeCh)
})
}

func (r *Reloader) run() {
var debounceTimer, retryTimer *time.Timer
var debounceC, retryC <-chan time.Time
stopTimer := func(timer *time.Timer) {
if timer != nil {
timer.Stop()
}
}
defer func() {
stopTimer(debounceTimer)
stopTimer(retryTimer)
}()

reloadAndMaybeRetry := func(isRetry bool) {
if isRetry {
r.cfg.Loggers.Debug("Retrying flag data load after earlier failure")
} else {
r.cfg.Loggers.Info("Reloading flag data after detecting a change")
}
// A pending retry is superseded by this reload: it either succeeded, or it failed
// and arms a fresh retry below.
stopTimer(retryTimer)
retryTimer, retryC = nil, nil
if !r.reload() && r.cfg.RetryDelay > 0 {
retryTimer = time.NewTimer(r.cfg.RetryDelay)
retryC = retryTimer.C
}
}

for {
select {
case <-r.closeCh:
return
case <-r.triggerCh:
if r.cfg.DebounceDelay <= 0 {
reloadAndMaybeRetry(false)
continue
}
if debounceTimer == nil {
debounceTimer = time.NewTimer(r.cfg.DebounceDelay)
debounceC = debounceTimer.C
} else {
// Stop-then-Reset without a drain is safe under this module's Go version:
// an expired timer's tick cannot be delivered after Reset. At worst a tick
// already consumed by the debounceC case races a new Trigger, producing one
// extra serialized reload that skip-unchanged absorbs.
stopTimer(debounceTimer)
debounceTimer.Reset(r.cfg.DebounceDelay)
}
case <-r.retryRequestCh:
// A synchronous ReloadNow failed; arm the retry timer without reloading again
// immediately. An already-armed retry keeps its earlier deadline.
if retryTimer == nil {
retryTimer = time.NewTimer(r.cfg.RetryDelay)
retryC = retryTimer.C
}
case <-debounceC:
debounceTimer, debounceC = nil, nil
reloadAndMaybeRetry(false)
case <-retryC:
retryTimer, retryC = nil, nil
reloadAndMaybeRetry(true)
}
}
}

// reload performs one full load of all configured files, returning true on success (including
// a skipped no-op application). The whole set is re-read on every reload: entries are combined
// across files in order, so a change to one file can alter which file wins for a key.
func (r *Reloader) reload() bool {
r.reloadMu.Lock()
defer r.reloadMu.Unlock()

// A trigger already queued when Close was called can still reach here; the select in
// run() does not prioritize closeCh over triggerCh.
if r.isClosed() {
return true

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Returning true on a skip is a little weird.

}

docs := make([]Document, 0, len(r.cfg.Paths))
hasher := sha256.New()
for _, path := range r.cfg.Paths {
rawData, err := os.ReadFile(path) //nolint:gosec // G304: ok to read file into variable
if err != nil {
return r.fail(&ReadError{Err: wrapReadError(err), Path: path})
}
_, _ = hasher.Write(rawData)
_, _ = hasher.Write([]byte{0})
doc, err := parseDocument(rawData)
if err != nil {
return r.fail(&ReadError{Err: err, Path: path})
}
docs = append(docs, doc)
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is there a way to do this in a streaming fashion so all file content doesn't have to come into memory at once for the hash?


merged, err := Merge(r.cfg.DuplicateKeysHandling, docs...)
if err != nil {
return r.fail(err)
}

// Close may have happened while the files were being read; deliver nothing in that case.
if r.isClosed() {
return true
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is this check necessary? Or just a statistically good spot to put it after some IO?


// A success right after a failure must apply even when the content is unchanged since
// the last success: the consumer heard about the failure through OnError and may have
// moved to an interrupted state, and only Apply tells it things are good again.
recovering := r.lastErrorMsg != ""
r.lastErrorMsg = ""
hash := hasher.Sum(nil)
if r.cfg.SkipUnchanged && !recovering && bytes.Equal(hash, r.lastGoodHash) {
return true
}
r.lastGoodHash = hash
if r.cfg.Apply != nil {
r.cfg.Apply(merged)
}
return true
Comment thread
kinyoklion marked this conversation as resolved.
}

func (r *Reloader) fail(err error) bool {
// Close may have happened while the files were being read; deliver nothing in that
// case, and report success so no retry is armed.
if r.isClosed() {
return true
}
// With automatic retries, a persistent failure would repeat the same log entry and the
// same callback on every attempt, so repeats of an identical failure are demoted to
// debug level and do not re-invoke OnError. Status-listening consumers therefore see
// one report per distinct failure rather than one per retry.
if err.Error() == r.lastErrorMsg {
r.cfg.Loggers.Debugf("Unable to load flags: %s", err)
return false
}
r.lastErrorMsg = err.Error()
r.cfg.Loggers.Errorf("Unable to load flags: %s", err)
if r.cfg.OnError != nil {
r.cfg.OnError(err)
}
return false
}

func (r *Reloader) isClosed() bool {
select {
case <-r.closeCh:
return true
default:
return false
}
}
Loading
Loading