codekingpro/portable-devtools
114k
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 