Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
discovery.py436 linesDownload Raw Back to dynamic
1# Copyright 2019 The Kubernetes Authors.2#3# Licensed under the Apache License, Version 2.0 (the "License");4# you may not use this file except in compliance with the License.5# You may obtain a copy of the License at6#7#     http://www.apache.org/licenses/LICENSE-2.08#9# Unless required by applicable law or agreed to in writing, software10# distributed under the License is distributed on an "AS IS" BASIS,11# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.12# See the License for the specific language governing permissions and13# limitations under the License.14 15import os16import six17import json18import logging19import hashlib20import tempfile21from functools import partial22from collections import defaultdict23from abc import abstractmethod, abstractproperty24 25from json.decoder import JSONDecodeError26from urllib3.exceptions import ProtocolError, MaxRetryError27 28from kubernetes import __version__29from .exceptions import NotFoundError, ResourceNotFoundError, ResourceNotUniqueError, ApiException, ServiceUnavailableError30from .resource import Resource, ResourceList31 32 33DISCOVERY_PREFIX = 'apis'34 35 36class Discoverer(object):37    """38        A convenient container for storing discovered API resources. Allows39        easy searching and retrieval of specific resources.40 41        Subclasses implement the abstract methods with different loading strategies.42    """43 44    def __init__(self, client, cache_file):45        self.client = client46        default_cache_id = self.client.configuration.host47        if six.PY3:48            default_cache_id = default_cache_id.encode('utf-8')49        try:50            default_cachefile_name = 'osrcp-{0}.json'.format(hashlib.md5(default_cache_id, usedforsecurity=False).hexdigest())51        except TypeError:52            # usedforsecurity is only supported in 3.9+53            default_cachefile_name = 'osrcp-{0}.json'.format(hashlib.md5(default_cache_id).hexdigest())54        self.__cache_file = cache_file or os.path.join(tempfile.gettempdir(), default_cachefile_name)55        self.__init_cache()56 57    def __init_cache(self, refresh=False):58        if refresh or not os.path.exists(self.__cache_file):59            self._cache = {'library_version': __version__}60            refresh = True61        else:62            try:63                with open(self.__cache_file, 'r') as f:64                    self._cache = json.load(f, cls=partial(CacheDecoder, self.client))65                if self._cache.get('library_version') != __version__:66                    # Version mismatch, need to refresh cache67                    self.invalidate_cache()68            except Exception as e:69                logging.error("load cache error: %s", e)70                self.invalidate_cache()71        self._load_server_info()72        self.discover()73        if refresh:74            self._write_cache()75 76    def _write_cache(self):77        try:78            with open(self.__cache_file, 'w') as f:79                json.dump(self._cache, f, cls=CacheEncoder)80        except Exception:81            # Failing to write the cache isn't a big enough error to crash on82            pass83 84    def invalidate_cache(self):85        self.__init_cache(refresh=True)86 87    @abstractproperty88    def api_groups(self):89        pass90 91    @abstractmethod92    def search(self, prefix=None, group=None, api_version=None, kind=None, **kwargs):93        pass94 95    @abstractmethod96    def discover(self):97        pass98 99    @property100    def version(self):101        return self.__version102 103    def default_groups(self, request_resources=False):104        groups = {}105        groups['api'] = { '': {106            'v1': (ResourceGroup( True, resources=self.get_resources_for_api_version('api', '', 'v1', True) )107                if request_resources else ResourceGroup(True))108        }}109 110        groups[DISCOVERY_PREFIX] = {'': {111            'v1': ResourceGroup(True, resources = {"List": [ResourceList(self.client)]})112        }}113        return groups114 115    def parse_api_groups(self, request_resources=False, update=False):116        """ Discovers all API groups present in the cluster """117        if not self._cache.get('resources') or update:118            self._cache['resources'] = self._cache.get('resources', {})119            groups_response = self.client.request('GET', '/{}'.format(DISCOVERY_PREFIX)).groups120 121            groups = self.default_groups(request_resources=request_resources)122 123            for group in groups_response:124                new_group = {}125                for version_raw in group['versions']:126                    version = version_raw['version']127                    resource_group = self._cache.get('resources', {}).get(DISCOVERY_PREFIX, {}).get(group['name'], {}).get(version)128                    preferred = version_raw == group['preferredVersion']129                    resources = resource_group.resources if resource_group else {}130                    if request_resources:131                        resources = self.get_resources_for_api_version(DISCOVERY_PREFIX, group['name'], version, preferred)132                    new_group[version] = ResourceGroup(preferred, resources=resources)133                groups[DISCOVERY_PREFIX][group['name']] = new_group134            self._cache['resources'].update(groups)135            self._write_cache()136 137        return self._cache['resources']138 139    def _load_server_info(self):140        def just_json(_, serialized):141            return serialized142 143        if not self._cache.get('version'):144            try:145                self._cache['version'] = {146                    'kubernetes': self.client.request('get', '/version', serializer=just_json)147                }148            except (ValueError, MaxRetryError) as e:149                if isinstance(e, MaxRetryError) and not isinstance(e.reason, ProtocolError):150                    raise151                if not self.client.configuration.host.startswith("https://"):152                    raise ValueError("Host value %s should start with https:// when talking to HTTPS endpoint" %153                                     self.client.configuration.host)154                else:155                    raise156 157        self.__version = self._cache['version']158 159    def get_resources_for_api_version(self, prefix, group, version, preferred):160        """ returns a dictionary of resources associated with provided (prefix, group, version)"""161 162        resources = defaultdict(list)163        subresources = {}164 165        path = '/'.join(filter(None, [prefix, group, version]))166        try:167            resources_response = self.client.request('GET', path).resources or []168        except (ServiceUnavailableError, JSONDecodeError):169            # Handle both service unavailable errors and JSON decode errors170            resources_response = []171 172        resources_raw = list(filter(lambda resource: '/' not in resource['name'], resources_response))173        subresources_raw = list(filter(lambda resource: '/' in resource['name'], resources_response))174        for subresource in subresources_raw:175            resource, name = subresource['name'].split('/', 1)176            if not subresources.get(resource):177                subresources[resource] = {}178            subresources[resource][name] = subresource179 180        for resource in resources_raw:181            # Prevent duplicate keys182            for key in ('prefix', 'group', 'api_version', 'client', 'preferred'):183                resource.pop(key, None)184 185            resourceobj = Resource(186                prefix=prefix,187                group=group,188                api_version=version,189                client=self.client,190                preferred=preferred,191                subresources=subresources.get(resource['name']),192                **resource193            )194            resources[resource['kind']].append(resourceobj)195 196            resource_list = ResourceList(self.client, group=group, api_version=version, base_kind=resource['kind'])197            resources[resource_list.kind].append(resource_list)198        return resources199 200    def get(self, **kwargs):201        """ Same as search, but will throw an error if there are multiple or no202            results. If there are multiple results and only one is an exact match203            on api_version, that resource will be returned.204        """205        results = self.search(**kwargs)206        # If there are multiple matches, prefer exact matches on api_version207        if len(results) > 1 and kwargs.get('api_version'):208            results = [209                result for result in results if result.group_version == kwargs['api_version']210            ]211        # If there are multiple matches, prefer non-List kinds212        if len(results) > 1 and not all([isinstance(x, ResourceList) for x in results]):213            results = [result for result in results if not isinstance(result, ResourceList)]214        if len(results) == 1:215            return results[0]216        elif not results:217            raise ResourceNotFoundError('No matches found for {}'.format(kwargs))218        else:219            raise ResourceNotUniqueError('Multiple matches found for {}: {}'.format(kwargs, results))220 221 222class LazyDiscoverer(Discoverer):223    """ A convenient container for storing discovered API resources. Allows224        easy searching and retrieval of specific resources.225 226        Resources for the cluster are loaded lazily.227    """228 229    def __init__(self, client, cache_file):230        Discoverer.__init__(self, client, cache_file)231        self.__update_cache = False232 233    def discover(self):234        self.__resources = self.parse_api_groups(request_resources=False)235 236    def __maybe_write_cache(self):237        if self.__update_cache:238            self._write_cache()239            self.__update_cache = False240 241    @property242    def api_groups(self):243        return self.parse_api_groups(request_resources=False, update=True)['apis'].keys()244 245    def search(self, **kwargs):246        # In first call, ignore ResourceNotFoundError and set default value for results247        try:248            results = self.__search(self.__build_search(**kwargs), self.__resources, [])249        except ResourceNotFoundError:250            results = []251        if not results:252            self.invalidate_cache()253            results = self.__search(self.__build_search(**kwargs), self.__resources, [])254        self.__maybe_write_cache()255        return results256 257    def __search(self,  parts, resources, reqParams):258        part = parts[0]259        if part != '*':260 261            resourcePart = resources.get(part)262            if not resourcePart:263                return []264            elif isinstance(resourcePart, ResourceGroup):265                if len(reqParams) != 2:266                    raise ValueError("prefix and group params should be present, have %s" % reqParams)267                # Check if we've requested resources for this group268                if not resourcePart.resources:269                    prefix, group, version = reqParams[0], reqParams[1], part270                    try:271                        resourcePart.resources = self.get_resources_for_api_version(272                            prefix, group, part, resourcePart.preferred)273                    except NotFoundError:274                        raise ResourceNotFoundError275 276                    self._cache['resources'][prefix][group][version] = resourcePart277                    self.__update_cache = True278                return self.__search(parts[1:], resourcePart.resources, reqParams)279            elif isinstance(resourcePart, dict):280                # In this case parts [0] will be a specified prefix, group, version281                # as we recurse282                return self.__search(parts[1:], resourcePart, reqParams + [part] )283            else:284                if parts[1] != '*' and isinstance(parts[1], dict):285                    for _resource in resourcePart:286                        for term, value in parts[1].items():287                            if getattr(_resource, term) == value:288                                return [_resource]289 290                    return []291                else:292                    return resourcePart293        else:294            matches = []295            for key in resources.keys():296                matches.extend(self.__search([key] + parts[1:], resources, reqParams))297            return matches298 299    def __build_search(self, prefix=None, group=None, api_version=None, kind=None, **kwargs):300        if not group and api_version and '/' in api_version:301            group, api_version = api_version.split('/')302 303        items = [prefix, group, api_version, kind, kwargs]304        return list(map(lambda x: x or '*', items))305 306    def __iter__(self):307        for prefix, groups in self.__resources.items():308            for group, versions in groups.items():309                for version, rg in versions.items():310                    # Request resources for this groupVersion if we haven't yet311                    if not rg.resources:312                        rg.resources = self.get_resources_for_api_version(313                            prefix, group, version, rg.preferred)314                        self._cache['resources'][prefix][group][version] = rg315                        self.__update_cache = True316                    for _, resource in six.iteritems(rg.resources):317                        yield resource318        self.__maybe_write_cache()319 320 321class EagerDiscoverer(Discoverer):322    """ A convenient container for storing discovered API resources. Allows323        easy searching and retrieval of specific resources.324 325        All resources are discovered for the cluster upon object instantiation.326    """327 328    def update(self, resources):329        self.__resources = resources330 331    def __init__(self, client, cache_file):332        Discoverer.__init__(self, client, cache_file)333 334    def discover(self):335        self.__resources = self.parse_api_groups(request_resources=True)336 337    @property338    def api_groups(self):339        """ list available api groups """340        return self.parse_api_groups(request_resources=True, update=True)['apis'].keys()341 342 343    def search(self, **kwargs):344        """ Takes keyword arguments and returns matching resources. The search345            will happen in the following order:346                prefix: The api prefix for a resource, ie, /api, /oapi, /apis. Can usually be ignored347                group: The api group of a resource. Will also be extracted from api_version if it is present there348                api_version: The api version of a resource349                kind: The kind of the resource350                arbitrary arguments (see below), in random order351 352            The arbitrary arguments can be any valid attribute for an Resource object353        """354        results = self.__search(self.__build_search(**kwargs), self.__resources)355        if not results:356            self.invalidate_cache()357            results = self.__search(self.__build_search(**kwargs), self.__resources)358        return results359 360    def __build_search(self, prefix=None, group=None, api_version=None, kind=None, **kwargs):361        if not group and api_version and '/' in api_version:362            group, api_version = api_version.split('/')363 364        items = [prefix, group, api_version, kind, kwargs]365        return list(map(lambda x: x or '*', items))366 367    def __search(self, parts, resources):368        part = parts[0]369        resourcePart = resources.get(part)370 371        if part != '*' and resourcePart:372            if isinstance(resourcePart, ResourceGroup):373                return self.__search(parts[1:], resourcePart.resources)374            elif isinstance(resourcePart, dict):375                return self.__search(parts[1:], resourcePart)376            else:377                if parts[1] != '*' and isinstance(parts[1], dict):378                    for _resource in resourcePart:379                        for term, value in parts[1].items():380                            if getattr(_resource, term) == value:381                                return [_resource]382                    return []383                else:384                    return resourcePart385        elif part == '*':386            matches = []387            for key in resources.keys():388                matches.extend(self.__search([key] + parts[1:], resources))389            return matches390        return []391 392    def __iter__(self):393        for _, groups in self.__resources.items():394            for _, versions in groups.items():395                for _, resources in versions.items():396                    for _, resource in resources.items():397                        yield resource398 399 400class ResourceGroup(object):401    """Helper class for Discoverer container"""402    def __init__(self, preferred, resources=None):403        self.preferred = preferred404        self.resources = resources or {}405 406    def to_dict(self):407        return {408            '_type': 'ResourceGroup',409            'preferred': self.preferred,410            'resources': self.resources,411        }412 413 414class CacheEncoder(json.JSONEncoder):415 416    def default(self, o):417        return o.to_dict()418 419 420class CacheDecoder(json.JSONDecoder):421    def __init__(self, client, *args, **kwargs):422        self.client = client423        json.JSONDecoder.__init__(self, object_hook=self.object_hook, *args, **kwargs)424 425    def object_hook(self, obj):426        if '_type' not in obj:427            return obj428        _type = obj.pop('_type')429        if _type == 'Resource':430            return Resource(client=self.client, **obj)431        elif _type == 'ResourceList':432            return ResourceList(self.client, **obj)433        elif _type == 'ResourceGroup':434            return ResourceGroup(obj['preferred'], resources=self.object_hook(obj['resources']))435        return obj436 
codekingpro/portable-devtools · Team Ai