Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
extended_operation.py226 linesDownload Raw Back to api_core
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 
codekingpro/portable-devtools · Team Ai