gcpdiag.queries.apis_utils

GCP API-related utility functions.
def execute_concurrently( api: Any, requests: List[Any], context: gcpdiag.models.Context) -> Iterator[Tuple[Any, Optional[Any], Optional[Exception]]]:
30def execute_concurrently(
31  api: Any, requests: List[Any], context: models.Context
32) -> Iterator[Tuple[Any, Optional[Any], Optional[Exception]]]:
33  """
34  Executes a list of API requests concurrently.
35  Uses ThreadPoolExecutor in API server context, batch_execute_all in CLI context.
36  Yields: (request, response, exception)
37  """
38  if not requests:
39    return
40
41  if context.context_provider:
42    # API Server context: Use ThreadPoolExecutor
43    exec_ = executor.get_executor(context)
44    future_to_request = {exec_.submit(execute_single_request, req): req for req in requests}
45
46    for future in concurrent.futures.as_completed(future_to_request):
47      request = future_to_request[future]
48      try:
49        response, exception = future.result()
50        yield (request, response, exception)
51      except googleapiclient.errors.HttpError as e:
52        yield (request, None, e)
53  else:
54    # CLI context: Use original batch_execute_all
55    yield from batch_execute_all(api, requests)

Executes a list of API requests concurrently. Uses ThreadPoolExecutor in API server context, batch_execute_all in CLI context. Yields: (request, response, exception)

def execute_concurrently_with_pagination( api: Any, requests: List[Any], next_function: Callable, context: gcpdiag.models.Context, log_text: Optional[str] = None, response_keyword: str = 'items') -> Iterator[Any]:
 92def execute_concurrently_with_pagination(
 93  api: Any,
 94  requests: List[Any],
 95  next_function: Callable,
 96  context: models.Context,
 97  log_text: Optional[str] = None,
 98  response_keyword: str = 'items',
 99) -> Iterator[Any]:
100  """
101  Executes and paginates a list of API 'list' requests concurrently.
102  """
103  if not context.context_provider:
104    yield from batch_list_all(api, requests, next_function, log_text, response_keyword)
105    return
106
107  yield from _execute_with_pagination_in_api_context(
108    api, requests, next_function, context, response_keyword
109  )

Executes and paginates a list of API 'list' requests concurrently.

def list_all( request, next_function: Callable, response_keyword='items') -> Iterator[Any]:
112def list_all(request, next_function: Callable, response_keyword='items') -> Iterator[Any]:
113  """Execute GCP API `request` and subsequently call `next_function` until
114  there are no more results. Assumes that it is a list method and that
115  the results are under a `items` key."""
116
117  while True:
118    try:
119      response = request.execute(num_retries=config.API_RETRIES)
120    except googleapiclient.errors.HttpError as err:
121      raise utils.GcpApiError(err) from err
122
123    # Empty lists are omitted in GCP API responses
124    if response_keyword in response:
125      yield from response[response_keyword]
126
127    request = next_function(previous_request=request, previous_response=response)
128    if request is None:
129      break

Execute GCP API request and subsequently call next_function until there are no more results. Assumes that it is a list method and that the results are under a items key.

def multi_list_all(requests: list, next_function: Callable) -> Iterator[Any]:
132def multi_list_all(
133  requests: list,
134  next_function: Callable,
135) -> Iterator[Any]:
136  for req in requests:
137    yield from list_all(req, next_function)
def batch_list_all( api, requests: list, next_function: Callable, log_text: Optional[str] = None, response_keyword='items'):
140def batch_list_all(
141  api,
142  requests: list,
143  next_function: Callable,
144  log_text: Optional[str] = None,
145  response_keyword='items',
146):
147  """Similar to list_all but using batch API except in TPC environment."""
148
149  if 'googleapis.com' not in requests[0].uri:
150    #  the api client library does not handle batch api calls for TPC yet, so
151    #  the batch is processed and collected one at a time in that case
152    for req in requests:
153      yield from list_all(req, next_function)
154  else:
155    yield from _original_batch(api, requests, next_function, log_text, response_keyword)

Similar to list_all but using batch API except in TPC environment.

def should_retry(resp_status):
195def should_retry(resp_status):
196  if resp_status >= 500:
197    return True
198  if resp_status == 429:  # too many requests
199    return True
200  return False
def get_nth_exponential_random_retry(n, random_pct, multiplier, random_fn=None):
203def get_nth_exponential_random_retry(n, random_pct, multiplier, random_fn=None):
204  random_fn = random_fn or random.random
205  return (1 - random_fn() * random_pct) * multiplier**n
def batch_execute_all(api, requests: list):
208def batch_execute_all(api, requests: list):
209  """Execute all `requests` using the batch API and yield (request,response,exception)
210  tuples."""
211  # results: (request, result, exception) tuples
212  results: List[Tuple[Any, Optional[Any], Optional[Exception]]] = []
213  requests_todo = requests
214  requests_in_flight: List = []
215  retry_count = 0
216
217  def fetch_all_cb(request_id, response, exception):
218    try:
219      request = requests_in_flight[int(request_id)]
220    except (IndexError, ValueError, TypeError):
221      logging.debug(
222        'BUG: Cannot find request %r in list of pending requests, dropping request.', request_id
223      )
224      return
225
226    if exception:
227      if (
228        isinstance(exception, googleapiclient.errors.HttpError)
229        and should_retry(exception.status_code)
230        and retry_count < config.API_RETRIES
231      ):
232        logging.debug(
233          'received HTTP error status code %d from API, retrying', exception.status_code
234        )
235        requests_todo.append(request)
236      else:
237        results.append((request, None, utils.GcpApiError(exception)))
238      return
239
240    if not response:
241      return
242
243    results.append((request, response, None))
244
245  while True:
246    requests_in_flight = requests_todo
247    requests_todo = []
248    results = []
249
250    # Do the batch API request
251    try:
252      batch = api.new_batch_http_request()
253      for i, req in enumerate(requests_in_flight):
254        batch.add(req, callback=fetch_all_cb, request_id=str(i))
255      batch.execute()
256    except (googleapiclient.errors.HttpError, httplib2.HttpLib2Error) as err:
257      if isinstance(err, googleapiclient.errors.HttpError):
258        error_msg = f'received HTTP error status code {err.status_code} from Batch API, retrying'
259      else:
260        error_msg = f'received exception from Batch API: {err}, retrying'
261      if (
262        not isinstance(err, googleapiclient.errors.HttpError) or should_retry(err.status_code)
263      ) and retry_count < config.API_RETRIES:
264        logging.debug(error_msg)
265        requests_todo = requests_in_flight
266        results = []
267      else:
268        raise utils.GcpApiError(err) from err
269
270    # Yield results
271    yield from results
272
273    # If no requests_todo, means we are done.
274    if not requests_todo:
275      break
276
277    # for example: retry delay: 20% is random, progression: 1, 1.4, 2.0, 2.7, ... 28.9 (10 retries)
278    sleep_time = get_nth_exponential_random_retry(
279      n=retry_count,
280      random_pct=config.API_RETRY_SLEEP_RANDOMNESS_PCT,
281      multiplier=config.API_RETRY_SLEEP_MULTIPLIER,
282    )
283    logging.debug('sleeping %.2f seconds before retry #%d', sleep_time, retry_count + 1)
284    time.sleep(sleep_time)
285    retry_count += 1

Execute all requests using the batch API and yield (request,response,exception) tuples.

def execute_single_request(request: Any) -> Tuple[Optional[Any], Optional[Exception]]:
288def execute_single_request(request: Any) -> Tuple[Optional[Any], Optional[Exception]]:
289  """Executes a single API request and returns the response and exception."""
290  try:
291    response = request.execute(num_retries=config.API_RETRIES)
292    return response, None
293  except googleapiclient.errors.HttpError as e:
294    return None, e

Executes a single API request and returns the response and exception.