|
| 1 | +""" |
| 2 | +This submodule contains the async in-memory feature store implementation. |
| 3 | +""" |
| 4 | + |
| 5 | +from collections import defaultdict |
| 6 | +from typing import Any, Dict, Mapping, Optional |
| 7 | + |
| 8 | +from ldclient.impl.util import log |
| 9 | +from ldclient.interfaces import AsyncFeatureStore, DiagnosticDescription |
| 10 | +from ldclient.versioned_data_kind import VersionedDataKind |
| 11 | + |
| 12 | + |
| 13 | +class AsyncInMemoryFeatureStore(AsyncFeatureStore, DiagnosticDescription): |
| 14 | + """The default async feature store implementation, which holds all data in memory. |
| 15 | +
|
| 16 | + .. caution:: |
| 17 | + This feature is experimental and should NOT be considered ready for production |
| 18 | + use. It may change or be removed without notice and is not subject to backwards |
| 19 | + compatibility guarantees. Pin to a specific minor version and review the changelog |
| 20 | + before upgrading. |
| 21 | + """ |
| 22 | + |
| 23 | + def __init__(self): |
| 24 | + """Constructs an instance of AsyncInMemoryFeatureStore.""" |
| 25 | + self._initialized = False |
| 26 | + self._items: Dict[VersionedDataKind, Dict[str, Any]] = defaultdict(dict) |
| 27 | + |
| 28 | + def is_monitoring_enabled(self) -> bool: |
| 29 | + return False |
| 30 | + |
| 31 | + def is_available(self) -> bool: |
| 32 | + return True |
| 33 | + |
| 34 | + async def get(self, kind: VersionedDataKind, key: str) -> Optional[Any]: |
| 35 | + """ """ |
| 36 | + items_of_kind = self._items[kind] |
| 37 | + item = items_of_kind.get(key) |
| 38 | + if item is None: |
| 39 | + log.debug("Attempted to get missing key %s in '%s', returning None", key, kind.namespace) |
| 40 | + return None |
| 41 | + if 'deleted' in item and item['deleted']: |
| 42 | + log.debug("Attempted to get deleted key %s in '%s', returning None", key, kind.namespace) |
| 43 | + return None |
| 44 | + return item |
| 45 | + |
| 46 | + async def all(self, kind: VersionedDataKind) -> Dict[str, Any]: |
| 47 | + """ """ |
| 48 | + items_of_kind = self._items[kind] |
| 49 | + return dict((k, i) for k, i in items_of_kind.items() if ('deleted' not in i) or not i['deleted']) |
| 50 | + |
| 51 | + async def init(self, all_data: Mapping[VersionedDataKind, Mapping[str, dict]]) -> None: |
| 52 | + """ """ |
| 53 | + all_decoded = {} |
| 54 | + for kind, items in all_data.items(): |
| 55 | + items_decoded = {} |
| 56 | + for key, item in items.items(): |
| 57 | + items_decoded[key] = kind.decode(item) |
| 58 | + all_decoded[kind] = items_decoded |
| 59 | + self._items.clear() |
| 60 | + self._items.update(all_decoded) |
| 61 | + self._initialized = True |
| 62 | + for k in all_data: |
| 63 | + log.debug("Initialized '%s' store with %d items", k.namespace, len(all_data[k])) |
| 64 | + |
| 65 | + async def delete(self, kind: VersionedDataKind, key: str, version: int) -> None: |
| 66 | + """ """ |
| 67 | + items_of_kind = self._items[kind] |
| 68 | + i = items_of_kind.get(key) |
| 69 | + if i is None or i['version'] < version: |
| 70 | + items_of_kind[key] = {'deleted': True, 'version': version} |
| 71 | + |
| 72 | + async def upsert(self, kind: VersionedDataKind, item: dict) -> bool: |
| 73 | + """ """ |
| 74 | + decoded_item = kind.decode(item) |
| 75 | + key = item['key'] |
| 76 | + items_of_kind = self._items[kind] |
| 77 | + i = items_of_kind.get(key) |
| 78 | + if i is None or i['version'] < item['version']: |
| 79 | + items_of_kind[key] = decoded_item |
| 80 | + log.debug("Updated %s in '%s' to version %d", key, kind.namespace, item['version']) |
| 81 | + return True |
| 82 | + return False |
| 83 | + |
| 84 | + @property |
| 85 | + def initialized(self) -> bool: |
| 86 | + """ """ |
| 87 | + return self._initialized |
| 88 | + |
| 89 | + async def close(self) -> None: |
| 90 | + """ """ |
| 91 | + pass |
| 92 | + |
| 93 | + def describe_configuration(self, config) -> str: |
| 94 | + return 'memory' |
0 commit comments