Skip to content

Commit f55dc74

Browse files
committed
feat: Add shared file data code with reload, retry, and polling
Adds ldclient.impl.integrations.files.filedata, the file reading, parsing, and merging logic shared by the components that load flag and segment data from local files. A document that starts with an opening brace is parsed as JSON and any other document as YAML, so JSON files no longer depend on the YAML parser. Definitions are decoded into the flag and segment models while the file is read, so an invalid definition fails that load instead of failing later inside the store. Files are merged in the configured order with a duplicate keys handling of fail or ignore. The Reloader owns the reload cycle: it serializes reloads, debounces change signals with a settle window, keeps the last good data by not applying a failed load, retries a failed load after a bounded delay, reports an identical failure once, and skips an application whose file contents did not change. It can treat a configured file that does not exist as a file with no content. The Poller detects changes by comparing modification time and size on an interval, including files that appear or disappear. The Watcher uses the watchdog package, watches the directory of each file so an absent file is picked up when it appears, matches the destination of a move so a file written by rename is detected, and retries a directory that does not exist yet. Both file data sources are rebuilt on this module. Their files must still exist and duplicate keys still fail the load. The FDv2 synchronizer now keeps the last good data across a failed load, reports an interrupted state, and retries, where before a failed initial load ended the synchronizer. The expansion of a flagValues entry is now an off flag serving its single variation, which is the form the Go SDK uses, so the evaluation reason kind for such a flag is OFF.
1 parent 1087147 commit f55dc74

5 files changed

Lines changed: 1664 additions & 467 deletions

File tree

Lines changed: 45 additions & 144 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,18 @@
1-
import json
2-
import os
31
import time
42
import traceback
5-
from typing import Optional
6-
7-
from ldclient.impl.repeating_task import RepeatingTask
3+
from typing import Optional, Union
4+
5+
from ldclient.impl.integrations.files.filedata import (
6+
DEFAULT_DEBOUNCE_DELAY,
7+
DEFAULT_RETRY_DELAY,
8+
DuplicateKeysHandling,
9+
MergeResult,
10+
Poller,
11+
Reloader,
12+
Watcher,
13+
abs_file_paths,
14+
have_watchdog
15+
)
816
from ldclient.impl.util import log
917
from ldclient.interfaces import (
1018
DataSourceErrorInfo,
@@ -15,43 +23,29 @@
1523
)
1624
from ldclient.versioned_data_kind import FEATURES, SEGMENTS
1725

18-
have_yaml = False
19-
try:
20-
import yaml
21-
22-
have_yaml = True
23-
except ImportError:
24-
pass
25-
26-
have_watchdog = False
27-
try:
28-
import watchdog
29-
import watchdog.events
30-
import watchdog.observers
31-
32-
have_watchdog = True
33-
except ImportError:
34-
pass
35-
36-
37-
def _sanitize_json_item(item):
38-
if not ('version' in item):
39-
item['version'] = 1
40-
4126

4227
class _FileDataSource(UpdateProcessor):
4328
def __init__(self, store, data_source_update_sink: Optional[DataSourceUpdateSink], ready, paths, auto_update, poll_interval, force_polling):
4429
self._store = store
4530
self._data_source_update_sink = data_source_update_sink
4631
self._ready = ready
4732
self._inited = False
48-
self._paths = paths
49-
if isinstance(self._paths, str):
50-
self._paths = [self._paths]
33+
self._paths = abs_file_paths(paths if isinstance(paths, list) else [paths])
5134
self._auto_update = auto_update
52-
self._auto_updater = None
35+
self._auto_updater: Optional[Union[Poller, Watcher]] = None
5336
self._poll_interval = poll_interval
5437
self._force_polling = force_polling
38+
# Debouncing and automatic retries only matter when something can trigger further
39+
# reloads. A source without auto update loads exactly once.
40+
self._reloader = Reloader(
41+
self._paths,
42+
DuplicateKeysHandling.FAIL,
43+
apply=self._apply,
44+
on_error=self._handle_error,
45+
debounce_delay=DEFAULT_DEBOUNCE_DELAY if auto_update else 0.0,
46+
retry_delay=DEFAULT_RETRY_DELAY if auto_update else 0.0,
47+
skip_unchanged=True,
48+
)
5549

5650
def _sink_or_store(self):
5751
"""
@@ -71,7 +65,7 @@ def _sink_or_store(self):
7165
return self._data_source_update_sink
7266

7367
def start(self):
74-
self._load_all()
68+
self._reloader.reload_now()
7569

7670
if self._auto_update:
7771
self._auto_updater = self._start_auto_updater()
@@ -81,23 +75,19 @@ def start(self):
8175
self._ready.set()
8276

8377
def stop(self):
84-
if self._auto_updater:
85-
self._auto_updater.stop()
78+
if self._auto_updater is not None:
79+
self._auto_updater.close()
80+
self._auto_updater = None
81+
self._reloader.close()
8682

8783
def initialized(self):
8884
return self._inited
8985

90-
def _load_all(self):
91-
all_data = {FEATURES: {}, SEGMENTS: {}}
92-
for path in self._paths:
93-
try:
94-
self._load_file(path, all_data)
95-
except Exception as e:
96-
log.error('Unable to load flag data from "%s": %s' % (path, repr(e)))
97-
traceback.print_exc()
98-
if self._data_source_update_sink is not None:
99-
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, DataSourceErrorInfo(DataSourceErrorKind.INVALID_DATA, 0, time.time, str(e)))
100-
return
86+
def _apply(self, merged: MergeResult):
87+
all_data = {
88+
FEATURES: {key: flag.to_json_dict() for key, flag in merged.flags.items()},
89+
SEGMENTS: {key: segment.to_json_dict() for key, segment in merged.segments.items()},
90+
}
10191
try:
10292
self._sink_or_store().init(all_data)
10393
self._inited = True
@@ -107,104 +97,15 @@ def _load_all(self):
10797
log.error('Unable to store data: %s' % repr(e))
10898
traceback.print_exc()
10999
if self._data_source_update_sink is not None:
110-
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, DataSourceErrorInfo(DataSourceErrorKind.UNKNOWN, 0, time.time, str(e)))
111-
112-
def _load_file(self, path, all_data):
113-
content = None
114-
with open(path, 'r') as f:
115-
content = f.read()
116-
parsed = self._parse_content(content)
117-
for key, flag in parsed.get('flags', {}).items():
118-
_sanitize_json_item(flag)
119-
self._add_item(all_data, FEATURES, flag)
120-
for key, value in parsed.get('flagValues', {}).items():
121-
self._add_item(all_data, FEATURES, self._make_flag_with_value(key, value))
122-
for key, segment in parsed.get('segments', {}).items():
123-
_sanitize_json_item(segment)
124-
self._add_item(all_data, SEGMENTS, segment)
125-
126-
def _parse_content(self, content):
127-
if have_yaml:
128-
return yaml.safe_load(content) # pyyaml correctly parses JSON too
129-
return json.loads(content)
100+
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, DataSourceErrorInfo(DataSourceErrorKind.UNKNOWN, 0, time.time(), str(e)))
130101

131-
def _add_item(self, all_data, kind, item):
132-
items = all_data[kind]
133-
key = item.get('key')
134-
if items.get(key) is None:
135-
items[key] = item
136-
else:
137-
raise Exception('In %s, key "%s" was used more than once' % (kind.namespace, key))
102+
def _handle_error(self, error: Exception):
103+
if self._data_source_update_sink is not None:
104+
self._data_source_update_sink.update_status(DataSourceState.INTERRUPTED, DataSourceErrorInfo(DataSourceErrorKind.INVALID_DATA, 0, time.time(), str(error)))
138105

139-
def _make_flag_with_value(self, key, value):
140-
return {'key': key, 'version': 1, 'on': True, 'fallthrough': {'variation': 0}, 'variations': [value]}
141-
142-
def _start_auto_updater(self):
143-
resolved_paths = []
144-
for path in self._paths:
145-
try:
146-
resolved_paths.append(os.path.realpath(path))
147-
except Exception:
148-
log.warning('Cannot watch for changes to data file "%s" because it is an invalid path' % path)
106+
def _start_auto_updater(self) -> Union[Poller, Watcher]:
149107
if have_watchdog and not self._force_polling:
150-
return _FileDataSource.WatchdogAutoUpdater(resolved_paths, self._load_all)
151-
else:
152-
return _FileDataSource.PollingAutoUpdater(resolved_paths, self._load_all, self._poll_interval)
153-
154-
# Watch for changes to data files using the watchdog package. This uses native OS filesystem notifications
155-
# if available for the current platform.
156-
class WatchdogAutoUpdater:
157-
def __init__(self, resolved_paths, reloader):
158-
watched_files = set(resolved_paths)
159-
160-
class LDWatchdogHandler(watchdog.events.FileSystemEventHandler):
161-
def on_any_event(self, event):
162-
if event.src_path in watched_files:
163-
reloader()
164-
165-
dir_paths = set()
166-
for path in resolved_paths:
167-
dir_paths.add(os.path.dirname(path))
168-
169-
self._observer = watchdog.observers.Observer()
170-
handler = LDWatchdogHandler()
171-
for path in dir_paths:
172-
self._observer.schedule(handler, path)
173-
self._observer.start()
174-
175-
def stop(self):
176-
self._observer.stop()
177-
self._observer.join()
178-
179-
# Watch for changes to data files by polling their modification times. This is used if auto-update is
180-
# on but the watchdog package is not installed.
181-
class PollingAutoUpdater:
182-
def __init__(self, resolved_paths, reloader, interval):
183-
self._paths = resolved_paths
184-
self._reloader = reloader
185-
self._file_times = self._check_file_times()
186-
self._timer = RepeatingTask.at_interval("ldclient.datasource.file.poll", interval, interval, self._poll)
187-
self._timer.start()
188-
189-
def stop(self):
190-
self._timer.stop()
191-
192-
def _poll(self):
193-
new_times = self._check_file_times()
194-
changed = False
195-
for file_path, file_time in self._file_times.items():
196-
if new_times.get(file_path) is not None and new_times.get(file_path) != file_time:
197-
changed = True
198-
break
199-
self._file_times = new_times
200-
if changed:
201-
self._reloader()
202-
203-
def _check_file_times(self):
204-
ret = {}
205-
for path in self._paths:
206-
try:
207-
ret[path] = os.path.getmtime(path)
208-
except Exception:
209-
ret[path] = None
210-
return ret
108+
return Watcher(self._paths, self._reloader.trigger)
109+
poller = Poller(self._paths, self._poll_interval, self._reloader.trigger)
110+
poller.start()
111+
return poller

0 commit comments

Comments
 (0)