Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
stream_util.py148 linesDownload Raw Back to foundation
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"""Helpful utilities related to the stream module."""
15
16import logging
17import threading
18
19from grpc.framework.foundation import stream
20
21_NO_VALUE = object()
22_LOGGER = logging.getLogger(__name__)
23
24
25class TransformingConsumer(stream.Consumer):
26    """A stream.Consumer that passes a transformation of its input to another."""
27
28    def __init__(self, transformation, downstream):
29        self._transformation = transformation
30        self._downstream = downstream
31
32    def consume(self, value):
33        self._downstream.consume(self._transformation(value))
34
35    def terminate(self):
36        self._downstream.terminate()
37
38    def consume_and_terminate(self, value):
39        self._downstream.consume_and_terminate(self._transformation(value))
40
41
42class IterableConsumer(stream.Consumer):
43    """A Consumer that when iterated over emits the values it has consumed."""
44
45    def __init__(self):
46        self._condition = threading.Condition()
47        self._values = []
48        self._active = True
49
50    def consume(self, value):
51        with self._condition:
52            if self._active:
53                self._values.append(value)
54                self._condition.notify()
55
56    def terminate(self):
57        with self._condition:
58            self._active = False
59            self._condition.notify()
60
61    def consume_and_terminate(self, value):
62        with self._condition:
63            if self._active:
64                self._values.append(value)
65                self._active = False
66                self._condition.notify()
67
68    def __iter__(self):
69        return self
70
71    def __next__(self):
72        return self.next()
73
74    def next(self):
75        with self._condition:
76            while self._active and not self._values:
77                self._condition.wait()
78            if self._values:
79                return self._values.pop(0)
80            raise StopIteration()
81
82
83class ThreadSwitchingConsumer(stream.Consumer):
84    """A Consumer decorator that affords serialization and asynchrony."""
85
86    def __init__(self, sink, pool):
87        self._lock = threading.Lock()
88        self._sink = sink
89        self._pool = pool
90        # True if self._spin has been submitted to the pool to be called once and
91        # that call has not yet returned, False otherwise.
92        self._spinning = False
93        self._values = []
94        self._active = True
95
96    def _spin(self, sink, value, terminate):
97        while True:
98            try:
99                if value is _NO_VALUE:
100                    sink.terminate()
101                elif terminate:
102                    sink.consume_and_terminate(value)
103                else:
104                    sink.consume(value)
105            except Exception as e:  # pylint:disable=broad-except
106                _LOGGER.exception(e)
107
108            with self._lock:
109                if terminate:
110                    self._spinning = False
111                    return
112                if self._values:
113                    value = self._values.pop(0)
114                    terminate = not self._values and not self._active
115                elif not self._active:
116                    value = _NO_VALUE
117                    terminate = True
118                else:
119                    self._spinning = False
120                    return
121
122    def consume(self, value):
123        with self._lock:
124            if self._active:
125                if self._spinning:
126                    self._values.append(value)
127                else:
128                    self._pool.submit(self._spin, self._sink, value, False)
129                    self._spinning = True
130
131    def terminate(self):
132        with self._lock:
133            if self._active:
134                self._active = False
135                if not self._spinning:
136                    self._pool.submit(self._spin, self._sink, _NO_VALUE, True)
137                    self._spinning = True
138
139    def consume_and_terminate(self, value):
140        with self._lock:
141            if self._active:
142                self._active = False
143                if self._spinning:
144                    self._values.append(value)
145                else:
146                    self._pool.submit(self._spin, self._sink, value, True)
147                    self._spinning = True
148 
codekingpro/portable-devtools · Team Ai