codekingpro/portable-devtools
114k
1# Copyright 2022 Google LLC2#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 15"""Futures for extended long-running operations returned from Google Cloud APIs.16 17These futures can be used to synchronously wait for the result of a18long-running operations using :meth:`ExtendedOperation.result`:19 20.. code-block:: python21 22 extended_operation = my_api_client.long_running_method()23 24 extended_operation.result()25 26Or asynchronously using callbacks and :meth:`Operation.add_done_callback`:27 28.. code-block:: python29 30 extended_operation = my_api_client.long_running_method()31 32 def my_callback(ex_op):33 print(f"Operation {ex_op.name} completed")34 35 extended_operation.add_done_callback(my_callback)36 37"""38 39import threading40 41from google.api_core import exceptions42from google.api_core.future import polling43 44 45class ExtendedOperation(polling.PollingFuture):46 """An ExtendedOperation future for interacting with a Google API Long-Running Operation.47 48 Args:49 extended_operation (proto.Message): The initial operation.50 refresh (Callable[[], type(extended_operation)]): A callable that returns51 the latest state of the operation.52 cancel (Callable[[], None]): A callable that tries to cancel the operation.53 polling Optional(google.api_core.retry.Retry): The configuration used54 for polling. This can be used to control how often :meth:`done`55 is polled. If the ``timeout`` argument to :meth:`result` is56 specified it will override the ``polling.timeout`` property.57 retry Optional(google.api_core.retry.Retry): DEPRECATED use ``polling``58 instead. If specified it will override ``polling`` parameter to59 maintain backward compatibility.60 61 Note: Most long-running API methods use google.api_core.operation.Operation62 This class is a wrapper for a subset of methods that use alternative63 Long-Running Operation (LRO) semantics.64 65 Note: there is not a concrete type the extended operation must be.66 It MUST have fields that correspond to the following, POSSIBLY WITH DIFFERENT NAMES:67 * name: str68 * status: Union[str, bool, enum.Enum]69 * error_code: int70 * error_message: str71 """72 73 def __init__(74 self,75 extended_operation,76 refresh,77 cancel,78 polling=polling.DEFAULT_POLLING,79 **kwargs,80 ):81 super().__init__(polling=polling, **kwargs)82 self._extended_operation = extended_operation83 self._refresh = refresh84 self._cancel = cancel85 # Note: the extended operation does not give a good way to indicate cancellation.86 # We make do with manually tracking cancellation and checking for doneness.87 self._cancelled = False88 self._completion_lock = threading.Lock()89 # Invoke in case the operation came back already complete.90 self._handle_refreshed_operation()91 92 # Note: the following four properties MUST be overridden in a subclass93 # if, and only if, the fields in the corresponding extended operation message94 # have different names.95 #96 # E.g. we have an extended operation class that looks like97 #98 # class MyOperation(proto.Message):99 # moniker = proto.Field(proto.STRING, number=1)100 # status_msg = proto.Field(proto.STRING, number=2)101 # optional http_error_code = proto.Field(proto.INT32, number=3)102 # optional http_error_msg = proto.Field(proto.STRING, number=4)103 #104 # the ExtendedOperation subclass would provide property overrides that map105 # to these (poorly named) fields.106 @property107 def name(self):108 return self._extended_operation.name109 110 @property111 def status(self):112 return self._extended_operation.status113 114 @property115 def error_code(self):116 return self._extended_operation.error_code117 118 @property119 def error_message(self):120 return self._extended_operation.error_message121 122 def __getattr__(self, name):123 return getattr(self._extended_operation, name)124 125 def done(self, retry=None):126 self._refresh_and_update(retry)127 return self._extended_operation.done128 129 def cancel(self):130 if self.done():131 return False132 133 self._cancel()134 self._cancelled = True135 return True136 137 def cancelled(self):138 # TODO(dovs): there is not currently a good way to determine whether the139 # operation has been cancelled.140 # The best we can do is manually keep track of cancellation141 # and check for doneness.142 if not self._cancelled:143 return False144 145 self._refresh_and_update()146 return self._extended_operation.done147 148 def _refresh_and_update(self, retry=None):149 if not self._extended_operation.done:150 self._extended_operation = (151 self._refresh(retry=retry) if retry else self._refresh()152 )153 self._handle_refreshed_operation()154 155 def _handle_refreshed_operation(self):156 with self._completion_lock:157 if not self._extended_operation.done:158 return159 160 if self.error_code and self.error_message:161 # Note: `errors` can be removed once proposal A from162 # b/284179390 is implemented.163 errors = []164 if hasattr(self, "error") and hasattr(self.error, "errors"):165 errors = self.error.errors166 exception = exceptions.from_http_status(167 status_code=self.error_code,168 message=self.error_message,169 response=self._extended_operation,170 errors=errors,171 )172 self.set_exception(exception)173 elif self.error_code or self.error_message:174 exception = exceptions.GoogleAPICallError(175 f"Unexpected error {self.error_code}: {self.error_message}"176 )177 self.set_exception(exception)178 else:179 # Extended operations have no payload.180 self.set_result(None)181 182 @classmethod183 def make(cls, refresh, cancel, extended_operation, **kwargs):184 """185 Return an instantiated ExtendedOperation (or child) that wraps186 * a refresh callable187 * a cancel callable (can be a no-op)188 * an initial result189 190 .. note::191 It is the caller's responsibility to set up refresh and cancel192 with their correct request argument.193 The reason for this is that the services that use Extended Operations194 have rpcs that look something like the following:195 196 // service.proto197 service MyLongService {198 rpc StartLongTask(StartLongTaskRequest) returns (ExtendedOperation) {199 option (google.cloud.operation_service) = "CustomOperationService";200 }201 }202 203 service CustomOperationService {204 rpc Get(GetOperationRequest) returns (ExtendedOperation) {205 option (google.cloud.operation_polling_method) = true;206 }207 }208 209 Any info needed for the poll, e.g. a name, path params, etc.210 is held in the request, which the initial client method is in a much211 better position to make made because the caller made the initial request.212 213 TL;DR: the caller sets up closures for refresh and cancel that carry214 the properly configured requests.215 216 Args:217 refresh (Callable[Optional[Retry]][type(extended_operation)]): A callable that218 returns the latest state of the operation.219 cancel (Callable[][Any]): A callable that tries to cancel the operation220 on a best effort basis.221 extended_operation (Any): The initial response of the long running method.222 See the docstring for ExtendedOperation.__init__ for requirements on223 the type and fields of extended_operation224 """225 return cls(extended_operation, refresh, cancel, **kwargs)226 