Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
utilities.py149 linesDownload Raw Back to beta
1# Copyright 2015 gRPC 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 at
6#
7#     http://www.apache.org/licenses/LICENSE-2.0
8#
9# Unless required by applicable law or agreed to in writing, software
10# 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 and
13# limitations under the License.
14"""Utilities for the gRPC Python Beta API."""
15
16import threading
17import time
18
19# implementations is referenced from specification in this module.
20from grpc.beta import implementations  # pylint: disable=unused-import
21from grpc.beta import interfaces
22from grpc.framework.foundation import callable_util
23from grpc.framework.foundation import future
24
25_DONE_CALLBACK_EXCEPTION_LOG_MESSAGE = (
26    'Exception calling connectivity future "done" callback!'
27)
28
29
30class _ChannelReadyFuture(future.Future):
31    def __init__(self, channel):
32        self._condition = threading.Condition()
33        self._channel = channel
34
35        self._matured = False
36        self._cancelled = False
37        self._done_callbacks = []
38
39    def _block(self, timeout):
40        until = None if timeout is None else time.time() + timeout
41        with self._condition:
42            while True:
43                if self._cancelled:
44                    raise future.CancelledError()
45                if self._matured:
46                    return
47                if until is None:
48                    self._condition.wait()
49                else:
50                    remaining = until - time.time()
51                    if remaining < 0:
52                        raise future.TimeoutError()
53                    self._condition.wait(timeout=remaining)
54
55    def _update(self, connectivity):
56        with self._condition:
57            if (
58                not self._cancelled
59                and connectivity is interfaces.ChannelConnectivity.READY
60            ):
61                self._matured = True
62                self._channel.unsubscribe(self._update)
63                self._condition.notify_all()
64                done_callbacks = tuple(self._done_callbacks)
65                self._done_callbacks = None
66            else:
67                return
68
69        for done_callback in done_callbacks:
70            callable_util.call_logging_exceptions(
71                done_callback, _DONE_CALLBACK_EXCEPTION_LOG_MESSAGE, self
72            )
73
74    def cancel(self):
75        with self._condition:
76            if not self._matured:
77                self._cancelled = True
78                self._channel.unsubscribe(self._update)
79                self._condition.notify_all()
80                done_callbacks = tuple(self._done_callbacks)
81                self._done_callbacks = None
82            else:
83                return False
84
85        for done_callback in done_callbacks:
86            callable_util.call_logging_exceptions(
87                done_callback, _DONE_CALLBACK_EXCEPTION_LOG_MESSAGE, self
88            )
89
90        return True
91
92    def cancelled(self):
93        with self._condition:
94            return self._cancelled
95
96    def running(self):
97        with self._condition:
98            return not self._cancelled and not self._matured
99
100    def done(self):
101        with self._condition:
102            return self._cancelled or self._matured
103
104    def result(self, timeout=None):
105        self._block(timeout)
106
107    def exception(self, timeout=None):
108        self._block(timeout)
109
110    def traceback(self, timeout=None):
111        self._block(timeout)
112
113    def add_done_callback(self, fn):
114        with self._condition:
115            if not self._cancelled and not self._matured:
116                self._done_callbacks.append(fn)
117                return
118
119        fn(self)
120
121    def start(self):
122        with self._condition:
123            self._channel.subscribe(self._update, try_to_connect=True)
124
125    def __del__(self):
126        with self._condition:
127            if not self._cancelled and not self._matured:
128                self._channel.unsubscribe(self._update)
129
130
131def channel_ready_future(channel):
132    """Creates a future.Future tracking when an implementations.Channel is ready.
133
134    Cancelling the returned future.Future does not tell the given
135    implementations.Channel to abandon attempts it may have been making to
136    connect; cancelling merely deactivates the return future.Future's
137    subscription to the given implementations.Channel's connectivity.
138
139    Args:
140      channel: An implementations.Channel.
141
142    Returns:
143      A future.Future that matures when the given Channel has connectivity
144        interfaces.ChannelConnectivity.READY.
145    """
146    ready_future = _ChannelReadyFuture(channel)
147    ready_future.start()
148    return ready_future
149 
codekingpro/portable-devtools · Team Ai