codekingpro/portable-devtools
114k
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 