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]:
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):
def
get_nth_exponential_random_retry(n, random_pct, multiplier, random_fn=None):
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.