|
4 | 4 |
|
5 | 5 | import threading |
6 | 6 | import traceback |
7 | | -from typing import Any, Callable, Dict, List, Mapping, Optional, Tuple |
| 7 | +from typing import Any, Callable, List, Optional, Tuple |
8 | 8 | from uuid import uuid4 |
9 | 9 |
|
10 | 10 | from ldclient.config import Config |
11 | 11 | from ldclient.context import Context |
12 | 12 | from ldclient.evaluation import EvaluationDetail, FeatureFlagsState |
13 | | -from ldclient.feature_store import _FeatureStoreDataSetSorter |
14 | 13 | from ldclient.hook import ( |
15 | 14 | EvaluationSeriesContext, |
16 | 15 | Hook, |
|
24 | 23 | from ldclient.impl.client_common import secure_mode_hash as _secure_mode_hash |
25 | 24 | from ldclient.impl.datasource.feature_requester import FeatureRequesterImpl |
26 | 25 | from ldclient.impl.datasource.polling import PollingUpdateProcessor |
27 | | -from ldclient.impl.datasource.status import ( |
28 | | - DataSourceStatusProviderImpl, |
29 | | - DataSourceUpdateSinkImpl |
30 | | -) |
31 | 26 | from ldclient.impl.datasource.streaming import StreamingUpdateProcessor |
32 | | -from ldclient.impl.datastore.status import ( |
33 | | - DataStoreStatusProviderImpl, |
34 | | - DataStoreUpdateSinkImpl |
35 | | -) |
36 | 27 | from ldclient.impl.datasystem import DataAvailability, DataSystem |
37 | 28 | from ldclient.impl.datasystem.fdv2 import FDv2 |
38 | 29 | from ldclient.impl.evaluator import Evaluator, error_reason |
|
43 | 34 | from ldclient.impl.events.event_processor import DefaultEventProcessor |
44 | 35 | from ldclient.impl.events.types import EventFactory |
45 | 36 | from ldclient.impl.flag_tracker import FlagTrackerImpl |
46 | | -from ldclient.impl.listeners import Listeners |
47 | 37 | from ldclient.impl.model.feature_flag import FeatureFlag |
48 | | -from ldclient.impl.repeating_task import RepeatingTask |
49 | 38 | from ldclient.impl.rwlock import ReadWriteLock |
50 | 39 | from ldclient.impl.stubs import NullEventProcessor, NullUpdateProcessor |
51 | 40 | from ldclient.impl.util import check_uwsgi, log |
52 | 41 | from ldclient.interfaces import ( |
53 | 42 | BigSegmentStoreStatusProvider, |
54 | 43 | DataSourceStatusProvider, |
55 | | - DataStoreStatus, |
56 | 44 | DataStoreStatusProvider, |
57 | | - DataStoreUpdateSink, |
58 | | - FeatureStore, |
59 | | - FlagTracker, |
60 | | - ReadOnlyStore |
| 45 | + FlagTracker |
61 | 46 | ) |
62 | 47 | from ldclient.migrations import OpTracker, Stage |
63 | 48 | from ldclient.plugin import EnvironmentMetadata |
64 | | -from ldclient.versioned_data_kind import FEATURES, SEGMENTS, VersionedDataKind |
| 49 | +from ldclient.versioned_data_kind import FEATURES, SEGMENTS |
65 | 50 |
|
66 | 51 | from .impl import AnyNum |
67 | 52 |
|
68 | 53 |
|
69 | | -class _FeatureStoreClientWrapper(FeatureStore): |
70 | | - """Provides additional behavior that the client requires before or after feature store operations. |
71 | | - Currently this just means sorting the data set for init() and dealing with data store status listeners. |
72 | | - """ |
73 | | - |
74 | | - def __init__(self, store: FeatureStore, store_update_sink: DataStoreUpdateSink): |
75 | | - self.store = store |
76 | | - self.__store_update_sink = store_update_sink |
77 | | - self.__monitoring_enabled = self.is_monitoring_enabled() |
78 | | - |
79 | | - # Covers the following variables |
80 | | - self.__lock = ReadWriteLock() |
81 | | - self.__last_available = True |
82 | | - self.__poller: Optional[RepeatingTask] = None |
83 | | - |
84 | | - def init(self, all_data: Mapping[VersionedDataKind, Mapping[str, Dict[Any, Any]]]): |
85 | | - return self.__wrapper(lambda: self.store.init(_FeatureStoreDataSetSorter.sort_all_collections(all_data))) |
86 | | - |
87 | | - def get(self, kind, key, callback): |
88 | | - return self.__wrapper(lambda: self.store.get(kind, key, callback)) |
89 | | - |
90 | | - def all(self, kind, callback): |
91 | | - return self.__wrapper(lambda: self.store.all(kind, callback)) |
92 | | - |
93 | | - def delete(self, kind, key, version): |
94 | | - return self.__wrapper(lambda: self.store.delete(kind, key, version)) |
95 | | - |
96 | | - def upsert(self, kind, item): |
97 | | - return self.__wrapper(lambda: self.store.upsert(kind, item)) |
98 | | - |
99 | | - @property |
100 | | - def initialized(self) -> bool: |
101 | | - return self.store.initialized |
102 | | - |
103 | | - def __wrapper(self, fn: Callable): |
104 | | - try: |
105 | | - return fn() |
106 | | - except BaseException: |
107 | | - if self.__monitoring_enabled: |
108 | | - self.__update_availability(False) |
109 | | - raise |
110 | | - |
111 | | - def __update_availability(self, available: bool): |
112 | | - with self.__lock.write(): |
113 | | - if available == self.__last_available: |
114 | | - return |
115 | | - self.__last_available = available |
116 | | - |
117 | | - status = DataStoreStatus(available, False) |
118 | | - |
119 | | - if available: |
120 | | - log.warn("Persistent store is available again") |
121 | | - |
122 | | - self.__store_update_sink.update_status(status) |
123 | | - |
124 | | - if available: |
125 | | - with self.__lock.write(): |
126 | | - if self.__poller is not None: |
127 | | - self.__poller.stop() |
128 | | - self.__poller = None |
129 | | - |
130 | | - return |
131 | | - |
132 | | - log.warn("Detected persistent store unavailability; updates will be cached until it recovers") |
133 | | - task = RepeatingTask("ldclient.check-availability", 0.5, 0, self.__check_availability) |
134 | | - |
135 | | - with self.__lock.write(): |
136 | | - self.__poller = task |
137 | | - self.__poller.start() |
138 | | - |
139 | | - def __check_availability(self): |
140 | | - try: |
141 | | - if self.store.is_available(): |
142 | | - self.__update_availability(True) |
143 | | - except BaseException as e: |
144 | | - log.error("Unexpected error from data store status function: %s", e) |
145 | | - |
146 | | - def is_monitoring_enabled(self) -> bool: |
147 | | - """ |
148 | | - This methods determines whether the wrapped store can support enabling monitoring. |
149 | | -
|
150 | | - The wrapped store must provide a monitoring_enabled method, which must |
151 | | - be true. But this alone is not sufficient. |
152 | | -
|
153 | | - Because this class wraps all interactions with a provided store, it can |
154 | | - technically "monitor" any store. However, monitoring also requires that |
155 | | - we notify listeners when the store is available again. |
156 | | -
|
157 | | - We determine this by checking the store's `available?` method, so this |
158 | | - is also a requirement for monitoring support. |
159 | | -
|
160 | | - These extra checks won't be necessary once `available` becomes a part |
161 | | - of the core interface requirements and this class no longer wraps every |
162 | | - feature store. |
163 | | - """ |
164 | | - |
165 | | - if not hasattr(self.store, 'is_monitoring_enabled'): |
166 | | - return False |
167 | | - |
168 | | - if not hasattr(self.store, 'is_available'): |
169 | | - return False |
170 | | - |
171 | | - monitoring_enabled = getattr(self.store, 'is_monitoring_enabled') |
172 | | - if not callable(monitoring_enabled): |
173 | | - return False |
174 | | - |
175 | | - return monitoring_enabled() |
176 | | - |
177 | | - |
178 | | -def _get_store_item(store, kind: VersionedDataKind, key: str) -> Any: |
179 | | - # This decorator around store.get provides backward compatibility with any custom data |
180 | | - # store implementation that might still be returning a dict, instead of our data model |
181 | | - # classes like FeatureFlag. |
182 | | - item = store.get(kind, key, lambda x: x) |
183 | | - return kind.decode(item) if isinstance(item, dict) else item |
184 | | - |
185 | | - |
186 | 54 | class LDClient: |
187 | 55 | """The LaunchDarkly SDK client object. |
188 | 56 |
|
@@ -268,8 +136,8 @@ def __start_up(self, start_wait: float): |
268 | 136 | self.__big_segment_store_manager = big_segment_store_manager |
269 | 137 |
|
270 | 138 | self._evaluator = Evaluator( |
271 | | - lambda key: _get_store_item(self._data_system.store, FEATURES, key), |
272 | | - lambda key: _get_store_item(self._data_system.store, SEGMENTS, key), |
| 139 | + lambda key: self._data_system.store.get(FEATURES, key), |
| 140 | + lambda key: self._data_system.store.get(SEGMENTS, key), |
273 | 141 | lambda key: big_segment_store_manager.get_user_membership(key), |
274 | 142 | log, |
275 | 143 | ) |
@@ -554,7 +422,7 @@ def _evaluate_internal(self, key: str, context: Context, default: Any, event_fac |
554 | 422 | return EvaluationDetail(default, None, error_reason('USER_NOT_SPECIFIED')), None |
555 | 423 |
|
556 | 424 | try: |
557 | | - flag = _get_store_item(self._data_system.store, FEATURES, key) |
| 425 | + flag = self._data_system.store.get(FEATURES, key) |
558 | 426 | except Exception as e: |
559 | 427 | log.error("Unexpected error while retrieving feature flag \"%s\": %s" % (key, repr(e))) |
560 | 428 | log.debug(traceback.format_exc()) |
|
0 commit comments