channel.py 31.8 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14
#   Copyright (c) 2020 PaddlePaddle Authors. All Rights Reserved.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
#     http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
# pylint: disable=doc-string-missing
B
barriery 已提交
15
from time import time as _time
D
dongdaxiang 已提交
16 17 18 19 20 21 22 23 24 25
import threading
import multiprocessing
import multiprocessing.queues
import sys
if sys.version_info.major == 2:
    import Queue
elif sys.version_info.major == 3:
    import queue as Queue
else:
    raise Exception("Error Python version")
26 27 28
import numpy as np
import logging
import enum
29
import os
30
import copy
D
dongdaxiang 已提交
31

32
_LOGGER = logging.getLogger(__name__)
B
barrierye 已提交
33

D
dongdaxiang 已提交
34 35 36 37 38 39 40

class ChannelDataEcode(enum.Enum):
    OK = 0
    TIMEOUT = 1
    NOT_IMPLEMENTED = 2
    TYPE_ERROR = 3
    RPC_PACKAGE_ERROR = 4
B
barrierye 已提交
41
    CLIENT_ERROR = 5
B
barrierye 已提交
42 43
    CLOSED_ERROR = 6
    UNKNOW = 7
D
dongdaxiang 已提交
44 45 46 47 48 49 50 51


class ChannelDataType(enum.Enum):
    DICT = 0
    CHANNEL_NPDATA = 1
    ERROR = 2


52 53 54 55
class ChannelData(object):
    def __init__(self,
                 datatype=None,
                 npdata=None,
B
barrierye 已提交
56
                 dictdata=None,
57 58
                 data_id=None,
                 ecode=None,
B
barrierye 已提交
59 60
                 error_info=None,
                 client_need_profile=False):
61 62 63
        '''
        There are several ways to use it:
        
B
barrierye 已提交
64 65 66
        1. ChannelData(ChannelDataType.CHANNEL_NPDATA.value, npdata, data_id)
        2. ChannelData(ChannelDataType.DICT.value, dictdata, data_id)
        3. ChannelData(ecode, error_info, data_id)
67 68 69 70 71 72

        Protobufs are not pickle-able:
        https://stackoverflow.com/questions/55344376/how-to-import-protobuf-module
        '''
        if ecode is not None:
            if data_id is None or error_info is None:
B
barriery 已提交
73 74
                _LOGGER.critical("Failed to generate ChannelData: data_id"
                                 " and error_info cannot be None")
75
                os._exit(-1)
76 77
            datatype = ChannelDataType.ERROR.value
        else:
B
barrierye 已提交
78 79
            if datatype == ChannelDataType.CHANNEL_NPDATA.value:
                ecode, error_info = ChannelData.check_npdata(npdata)
80
                if ecode != ChannelDataEcode.OK.value:
B
barrierye 已提交
81
                    datatype = ChannelDataType.ERROR.value
B
barriery 已提交
82
                    _LOGGER.error("(logid={}) {}".format(data_id, error_info))
B
barrierye 已提交
83 84 85 86
            elif datatype == ChannelDataType.DICT.value:
                ecode, error_info = ChannelData.check_dictdata(dictdata)
                if ecode != ChannelDataEcode.OK.value:
                    datatype = ChannelDataType.ERROR.value
B
barriery 已提交
87
                    _LOGGER.error("(logid={}) {}".format(data_id, error_info))
88
            else:
B
barriery 已提交
89 90
                _LOGGER.critical("(logid={}) datatype not match".format(
                    data_id))
91
                os._exit(-1)
92
        self.datatype = datatype
B
barrierye 已提交
93 94
        self.npdata = npdata
        self.dictdata = dictdata
95 96 97
        self.id = data_id
        self.ecode = ecode
        self.error_info = error_info
B
barrierye 已提交
98
        self.client_need_profile = client_need_profile
B
barrierye 已提交
99
        self.profile_data_set = set()
B
barrierye 已提交
100

B
barrierye 已提交
101
    def add_profile(self, profile_set):
B
barrierye 已提交
102 103
        if self.client_need_profile is False:
            self.client_need_profile = True
B
barrierye 已提交
104
        self.profile_data_set |= profile_set
105

B
barrierye 已提交
106 107 108 109
    @staticmethod
    def check_dictdata(dictdata):
        ecode = ChannelDataEcode.OK.value
        error_info = None
B
barrierye 已提交
110 111 112 113 114
        if isinstance(dictdata, list):
            # batch data
            for sample in dictdata:
                if not isinstance(sample, dict):
                    ecode = ChannelDataEcode.TYPE_ERROR.value
B
barriery 已提交
115 116
                    error_info = "Failed to check data: the type of " \
                            "data must be dict, but get {}.".format(type(sample))
B
barrierye 已提交
117 118 119
                    break
        elif not isinstance(dictdata, dict):
            # batch size = 1
B
barrierye 已提交
120
            ecode = ChannelDataEcode.TYPE_ERROR.value
B
barriery 已提交
121 122
            error_info = "Failed to check data: the type of data must " \
                    "be dict, but get {}.".format(type(dictdata))
B
barrierye 已提交
123
        return ecode, error_info
B
barrierye 已提交
124

B
bug fix  
barriery 已提交
125 126 127 128 129 130 131 132 133 134
    @staticmethod
    def check_batch_npdata(batch):
        ecode = ChannelDataEcode.OK.value
        error_info = None
        for npdata in batch:
            ecode, error_info = ChannelData.check_npdata(npdata)
            if ecode != ChannelDataEcode.OK.value:
                break
        return ecode, error_info

B
barrierye 已提交
135 136
    @staticmethod
    def check_npdata(npdata):
137 138
        ecode = ChannelDataEcode.OK.value
        error_info = None
W
wangjiawei04 已提交
139 140 141 142 143
        if isinstance(npdata, list):
            # batch data
            for sample in npdata:
                if not isinstance(sample, dict):
                    ecode = ChannelDataEcode.TYPE_ERROR.value
B
barriery 已提交
144 145 146
                    error_info = "Failed to check data: the " \
                            "value of data must be dict, but get {}.".format(
                                    type(sample))
W
wangjiawei04 已提交
147 148 149 150
                    break
                for _, value in sample.items():
                    if not isinstance(value, np.ndarray):
                        ecode = ChannelDataEcode.TYPE_ERROR.value
B
barriery 已提交
151 152 153
                        error_info = "Failed to check data: the" \
                                " value of data must be np.ndarray, but get {}.".format(
                                        type(value))
W
wangjiawei04 已提交
154 155 156 157 158 159
                        return ecode, error_info
        elif isinstance(npdata, dict):
            # batch_size = 1
            for _, value in npdata.items():
                if not isinstance(value, np.ndarray):
                    ecode = ChannelDataEcode.TYPE_ERROR.value
B
barriery 已提交
160 161 162
                    error_info = "Failed to check data: the value " \
                            "of data must be np.ndarray, but get {}.".format(
                                    type(value))
W
wangjiawei04 已提交
163 164 165
                    break
        else:
            ecode = ChannelDataEcode.TYPE_ERROR.value
B
barriery 已提交
166 167
            error_info = "Failed to check data: the value of data " \
                    "must be dict, but get {}.".format(type(npdata))
168 169 170 171
        return ecode, error_info

    def parse(self):
        feed = None
B
barrierye 已提交
172 173
        if self.datatype == ChannelDataType.CHANNEL_NPDATA.value:
            # return narray
174
            feed = self.npdata
B
barrierye 已提交
175 176 177
        elif self.datatype == ChannelDataType.DICT.value:
            # return dict
            feed = self.dictdata
178
        else:
B
barriery 已提交
179 180
            _LOGGER.critical("Failed to parse channeldata: error " \
                    "type({}) in datatype.".format(self.datatype))
181
            os._exit(-1)
182 183
        return feed

184 185 186 187 188 189 190 191
    def __cmp__(self, other):
        if self.id < other.id:
            return -1
        elif self.id == other.id:
            return 0
        else:
            return 1

192 193 194 195 196
    def __str__(self):
        return "type[{}], ecode[{}], id[{}]".format(
            ChannelDataType(self.datatype).name, self.ecode, self.id)


B
barrierye 已提交
197
class ProcessChannel(object):
198 199 200 201 202 203 204 205 206
    """ 
    (Process version) The channel used for communication between Ops.

    1. Support multiple different Op feed data (multiple producer)
        Different types of data will be packaged through the data ID
    2. Support multiple different Op fetch data (multiple consumer)
        Only when all types of Ops get the data of the same ID,
        the data will be poped; The Op of the same type will not
        get the data of the same ID.
B
barriery 已提交
207
    3. Function front support timeout param to make auto-batching.
208 209 210 211 212

    Note:
    1. The ID of the data in the channel must be different.
    2. The function add_producer() and add_consumer() are not thread safe,
       and can only be called during initialization.
B
barrierye 已提交
213 214 215 216 217 218 219 220 221 222 223

    There are two buffers and one queue in Channel:

        op_A \                                           / op_D
        op_B - a. input_buf -> b. queue -> c. output_buf - op_E
        op_C /                                           \ op_F
    
    a. In input_buf, the input of multiple predecessor Ops is packed by data ID.
    b. The packed data will be stored in queue.
    c. In order to support multiple successor Ops to retrieve data, output_buf
        maintains the data obtained from queue.
224 225
    """

B
barriery 已提交
226
    def __init__(self, manager, name=None, maxsize=0):
B
barrierye 已提交
227 228 229 230 231 232
        # For queue multiprocess: after putting an object on 
        # an empty queue there may be an infinitessimal delay
        # before the queue's :meth:`~Queue.empty`
        # see more:
        # - https://bugs.python.org/issue18277
        # - https://hg.python.org/cpython/rev/860fc6a2bd21
233
        self._que = manager.PriorityQueue(maxsize=maxsize)
234 235
        self._maxsize = maxsize
        self.name = name
236
        self._stop = manager.Value('i', 0)
237 238 239 240

        self._cv = multiprocessing.Condition()

        self._producers = []
B
barrierye 已提交
241
        self._pushed_producer_count = manager.dict()  # {data_id: count}
B
barrierye 已提交
242
        self._input_buf = manager.dict()  # {data_id: {op_name: data}}
243

244
        self._reset_max_cursor = 1000000000000000000
B
barrierye 已提交
245 246 247 248
        self._consumer_cursors = manager.dict()  # {op_name: cursor}
        self._cursor_count = manager.dict()  # {cursor: count}
        self._base_cursor = manager.Value('i', 0)
        self._output_buf = manager.list()
249

B
barriery 已提交
250 251 252
    def get_maxsize(self):
        return self._maxsize

B
barriery 已提交
253 254 255
    def size(self):
        return self._que.qsize()

256 257 258 259
    def get_producers(self):
        return self._producers

    def get_consumers(self):
B
barrierye 已提交
260
        return self._consumer_cursors.keys()
261 262 263 264 265 266 267

    def _log(self, info_str):
        return "[{}] {}".format(self.name, info_str)

    def add_producer(self, op_name):
        """ not thread safe, and can only be called during initialization. """
        if op_name in self._producers:
268
            _LOGGER.critical(
B
barriery 已提交
269 270
                self._log("Failed to add producer: producer({})" \
                        " is already in channel".format(op_name)))
271
            os._exit(-1)
272
        self._producers.append(op_name)
B
barriery 已提交
273
        _LOGGER.debug(self._log("Succ add a producer: {}".format(op_name)))
274 275 276

    def add_consumer(self, op_name):
        """ not thread safe, and can only be called during initialization. """
B
barrierye 已提交
277
        if op_name in self._consumer_cursors:
278
            _LOGGER.critical(
B
barriery 已提交
279 280
                    self._log("Failed to add consumer: consumer({})" \
                            " is already in channel".format(op_name)))
281
            os._exit(-1)
B
barrierye 已提交
282
        self._consumer_cursors[op_name] = 0
283

B
barrierye 已提交
284 285 286
        if self._cursor_count.get(0) is None:
            self._cursor_count[0] = 0
        self._cursor_count[0] += 1
B
barriery 已提交
287
        _LOGGER.debug(self._log("Succ add a consumer: {}".format(op_name)))
288 289

    def push(self, channeldata, op_name=None):
B
barrierye 已提交
290
        _LOGGER.debug(
B
barriery 已提交
291 292
            self._log("(logid={}) Op({}) Pushing data".format(channeldata.id,
                                                              op_name)))
293
        if len(self._producers) == 0:
294
            _LOGGER.critical(
295
                self._log(
B
barriery 已提交
296 297 298
                    "(logid={}) Op({}) Failed to push data: expected number"
                    " of producers to be greater than 0, but the it is 0.".
                    format(channeldata.id, op_name)))
299
            os._exit(-1)
300 301
        elif len(self._producers) == 1:
            with self._cv:
302
                while self._stop.value == 0:
303
                    try:
B
barrierye 已提交
304
                        self._que.put({op_name: channeldata}, timeout=0)
305 306 307
                        break
                    except Queue.Full:
                        self._cv.wait()
308
                if self._stop.value == 1:
B
barrierye 已提交
309
                    raise ChannelStopError()
310
                self._cv.notify_all()
B
barriery 已提交
311
            _LOGGER.debug(
B
barriery 已提交
312 313
                self._log("(logid={}) Op({}) Pushed data into internal queue.".
                          format(channeldata.id, op_name)))
314 315
            return True
        elif op_name is None:
316
            _LOGGER.critical(
317
                self._log(
B
barriery 已提交
318 319 320
                    "(logid={}) Op({}) Failed to push data: there are multiple "
                    "producers, so op_name cannot be None.".format(
                        channeldata.id, op_name)))
321
            os._exit(-1)
322 323 324 325 326

        producer_num = len(self._producers)
        data_id = channeldata.id
        put_data = None
        with self._cv:
B
barrierye 已提交
327 328
            if data_id not in self._input_buf:
                self._input_buf[data_id] = {
329 330 331
                    name: None
                    for name in self._producers
                }
B
barrierye 已提交
332
                self._pushed_producer_count[data_id] = 0
333
            # see: https://docs.python.org/3.6/library/multiprocessing.html?highlight=multiprocess#proxy-objects
B
barrierye 已提交
334 335 336 337 338
            # self._input_buf[data_id][op_name] = channeldata
            tmp_input_buf = self._input_buf[data_id]
            tmp_input_buf[op_name] = channeldata
            self._input_buf[data_id] = tmp_input_buf

B
barrierye 已提交
339
            if self._pushed_producer_count[data_id] + 1 == producer_num:
B
barrierye 已提交
340 341
                put_data = self._input_buf[data_id]
                self._input_buf.pop(data_id)
B
barrierye 已提交
342
                self._pushed_producer_count.pop(data_id)
343
            else:
B
barrierye 已提交
344
                self._pushed_producer_count[data_id] += 1
345 346

            if put_data is None:
B
barrierye 已提交
347
                _LOGGER.debug(
B
barriery 已提交
348 349 350
                    self._log(
                        "(logid={}) Op({}) Pushed data into input_buffer.".
                        format(data_id, op_name)))
351
            else:
352
                while self._stop.value == 0:
353
                    try:
B
barrierye 已提交
354
                        self._que.put(put_data, timeout=0)
355 356 357
                        break
                    except Queue.Empty:
                        self._cv.wait()
358
                if self._stop.value == 1:
B
barrierye 已提交
359
                    raise ChannelStopError()
360

B
barrierye 已提交
361
                _LOGGER.debug(
B
barriery 已提交
362 363 364
                    self._log(
                        "(logid={}) Op({}) Pushed data into internal_queue.".
                        format(data_id, op_name)))
365 366 367
            self._cv.notify_all()
        return True

B
barriery 已提交
368
    def front(self, op_name=None, timeout=None):
B
barriery 已提交
369
        _LOGGER.debug(
B
barriery 已提交
370 371
            self._log("Op({}) Getting data[?]; timeout(s)={}".format(op_name,
                                                                     timeout)))
B
barriery 已提交
372
        endtime = None
B
bug fix  
barriery 已提交
373 374 375 376 377
        if timeout is not None:
            if timeout <= 0:
                timeout = None
            else:
                endtime = _time() + timeout
B
barriery 已提交
378

B
barrierye 已提交
379
        if len(self._consumer_cursors) == 0:
380
            _LOGGER.critical(
381
                self._log(
B
barriery 已提交
382 383
                    "Op({}) Failed to get data: expected number of consumers to be " \
                            "greater than 0, but the it is 0.".format(op_name)))
384
            os._exit(-1)
B
barrierye 已提交
385
        elif len(self._consumer_cursors) == 1:
386 387
            resp = None
            with self._cv:
388
                while self._stop.value == 0 and resp is None:
389
                    try:
B
barrierye 已提交
390
                        resp = self._que.get(timeout=0)
391 392
                        break
                    except Queue.Empty:
B
barriery 已提交
393 394 395
                        if timeout is not None:
                            remaining = endtime - _time()
                            if remaining <= 0.0:
B
barriery 已提交
396
                                _LOGGER.debug(
B
barriery 已提交
397 398
                                    self._log("Op({}) Failed to get data: "
                                              "timeout".format(op_name)))
B
barriery 已提交
399 400 401 402
                                raise ChannelTimeoutError()
                            self._cv.wait(remaining)
                        else:
                            self._cv.wait()
403
                if self._stop.value == 1:
B
barrierye 已提交
404
                    raise ChannelStopError()
B
barriery 已提交
405
            _LOGGER.debug(
B
barriery 已提交
406 407
                self._log("(logid={}) Op({}) Got data".format(resp.values()[0]
                                                              .id, op_name)))
408 409
            return resp
        elif op_name is None:
410
            _LOGGER.critical(
411
                self._log(
B
barriery 已提交
412 413
                    "Op({}) Failed to get data: there are multiple consumers, "
                    "so op_name cannot be None.".format(op_name)))
414
            os._exit(-1)
415

B
barrierye 已提交
416 417 418 419 420 421 422 423 424 425
        # In output_buf, different Ops (according to op_name) have different
        # cursors. In addition, there is a base_cursor. Their difference is
        # the data_idx to be taken by the corresponding Op at the current
        # time:    data_idx = consumer_cursor - base_cursor
        # 
        #            base_cursor    consumer_B_cursor (data_idx: 3)
        #                 |                       |
        # output_buf: | data0 | data1 | data2 | data3 |
        #                 |
        #   consumer_A_cursor (data_idx: 0)
426
        with self._cv:
B
barrierye 已提交
427 428
            # When the data required by the current Op is not in output_buf,
            # it is necessary to obtain a data from queue and add it to output_buf.
429
            while self._stop.value == 0 and self._consumer_cursors[
B
barrierye 已提交
430
                    op_name] - self._base_cursor.value >= len(self._output_buf):
431
                try:
B
barrierye 已提交
432
                    channeldata = self._que.get(timeout=0)
B
barrierye 已提交
433
                    self._output_buf.append(channeldata)
B
barriery 已提交
434
                    _LOGGER.debug(
B
barriery 已提交
435 436 437
                        self._log(
                            "(logid={}) Op({}) Pop ready item into output_buffer".
                            format(channeldata.values()[0].id, op_name)))
438 439
                    break
                except Queue.Empty:
B
barriery 已提交
440 441 442
                    if timeout is not None:
                        remaining = endtime - _time()
                        if remaining <= 0.0:
B
barriery 已提交
443
                            _LOGGER.debug(
B
barriery 已提交
444 445
                                self._log("Op({}) Failed to get data: timeout".
                                          format(op_name)))
B
barriery 已提交
446 447 448 449
                            raise ChannelTimeoutError()
                        self._cv.wait(remaining)
                    else:
                        self._cv.wait()
450
            if self._stop.value == 1:
B
barrierye 已提交
451
                raise ChannelStopError()
452

B
barrierye 已提交
453 454 455 456
            consumer_cursor = self._consumer_cursors[op_name]
            base_cursor = self._base_cursor.value
            data_idx = consumer_cursor - base_cursor
            resp = self._output_buf[data_idx]
457

B
barrierye 已提交
458 459 460 461 462 463 464 465
            self._cursor_count[consumer_cursor] -= 1
            if consumer_cursor == base_cursor and self._cursor_count[
                    consumer_cursor] == 0:
                # When all the different Ops get the data that data_idx points
                # to, pop the data from output_buf.
                self._cursor_count.pop(consumer_cursor)
                self._output_buf.pop(0)
                self._base_cursor.value += 1
466 467
                # to avoid cursor overflow
                if self._base_cursor.value >= self._reset_max_cursor:
B
barriery 已提交
468
                    _LOGGER.info(self._log("Reset cursor in Channel"))
469 470 471 472 473 474 475 476 477 478
                    self._base_cursor.value -= self._reset_max_cursor
                    for name in self._consumer_cursors.keys():
                        self._consumer_cursors[name] -= self._reset_max_cursor
                    cursor_count_tmp = {
                        cursor - self._reset_max_cursor: count
                        for cursor, count in self._cursor_count.copy().items()
                    }
                    self._cursor_count.clear()
                    for cursor, count in cursor_count_tmp.items():
                        self._cursor_count[cursor] = count
B
barrierye 已提交
479 480 481 482 483 484

            self._consumer_cursors[op_name] += 1
            new_consumer_cursor = self._consumer_cursors[op_name]
            if self._cursor_count.get(new_consumer_cursor) is None:
                self._cursor_count[new_consumer_cursor] = 0
            self._cursor_count[new_consumer_cursor] += 1
485

486 487
            self._cv.notify_all()

B
barriery 已提交
488
        _LOGGER.debug(
B
barriery 已提交
489 490
            self._log("(logid={}) Op({}) Got data from output_buffer".format(
                resp.values()[0].id, op_name)))
B
barriery 已提交
491
        return resp
492 493

    def stop(self):
494
        _LOGGER.info(self._log("stop."))
495
        self._stop.value = 1
B
barrierye 已提交
496 497
        with self._cv:
            self._cv.notify_all()
498 499


500
class ThreadChannel(Queue.PriorityQueue):
501 502 503 504 505 506 507 508 509
    """ 
    (Thread version)The channel used for communication between Ops.

    1. Support multiple different Op feed data (multiple producer)
        Different types of data will be packaged through the data ID
    2. Support multiple different Op fetch data (multiple consumer)
        Only when all types of Ops get the data of the same ID,
        the data will be poped; The Op of the same type will not
        get the data of the same ID.
B
barriery 已提交
510
    3. Function front support timeout param to make auto-batching.
511 512 513 514 515

    Note:
    1. The ID of the data in the channel must be different.
    2. The function add_producer() and add_consumer() are not thread safe,
       and can only be called during initialization.
B
barrierye 已提交
516 517 518 519 520 521 522 523 524 525 526

    There are two buffers and one queue in Channel:

        op_A \                                           / op_D
        op_B - a. input_buf -> b. queue -> c. output_buf - op_E
        op_C /                                           \ op_F
    
    a. In input_buf, the input of multiple predecessor Ops is packed by data ID.
    b. The packed data will be stored in queue.
    c. In order to support multiple successor Ops to retrieve data, output_buf
        maintains the data obtained from queue.
527 528
    """

B
barriery 已提交
529
    def __init__(self, name=None, maxsize=-1):
530 531 532 533 534 535 536 537
        Queue.Queue.__init__(self, maxsize=maxsize)
        self._maxsize = maxsize
        self.name = name
        self._stop = False

        self._cv = threading.Condition()

        self._producers = []
B
barrierye 已提交
538
        self._pushed_producer_count = {}  # {data_id: count}
B
barrierye 已提交
539
        self._input_buf = {}  # {data_id: {op_name: data}}
540

541
        self._reset_max_cursor = 1000000000000000000
B
barrierye 已提交
542 543 544 545
        self._consumer_cursors = {}  # {op_name: idx}
        self._cursor_count = {}  # {cursor: count}
        self._base_cursor = 0
        self._output_buf = []
546

B
barriery 已提交
547 548 549
    def get_maxsize(self):
        return self._maxsize

B
barriery 已提交
550 551 552
    def size(self):
        return self.qsize()

553 554 555 556
    def get_producers(self):
        return self._producers

    def get_consumers(self):
B
barrierye 已提交
557
        return self._consumer_cursors.keys()
558 559 560 561 562 563 564

    def _log(self, info_str):
        return "[{}] {}".format(self.name, info_str)

    def add_producer(self, op_name):
        """ not thread safe, and can only be called during initialization. """
        if op_name in self._producers:
565
            _LOGGER.critical(
B
barriery 已提交
566 567
                self._log("Failed to add producer: producer({}) is "
                          "already in channel".format(op_name)))
568
            os._exit(-1)
569
        self._producers.append(op_name)
B
barriery 已提交
570
        _LOGGER.debug(self._log("Succ add a producer: {}".format(op_name)))
571 572 573

    def add_consumer(self, op_name):
        """ not thread safe, and can only be called during initialization. """
B
barrierye 已提交
574
        if op_name in self._consumer_cursors:
575
            _LOGGER.critical(
B
barriery 已提交
576 577
                self._log("Failed to add consumer: consumer({}) is "
                          "already in channel".format(op_name)))
578
            os._exit(-1)
B
barrierye 已提交
579
        self._consumer_cursors[op_name] = 0
580

B
barrierye 已提交
581 582 583
        if self._cursor_count.get(0) is None:
            self._cursor_count[0] = 0
        self._cursor_count[0] += 1
B
barriery 已提交
584
        _LOGGER.debug(self._log("Succ add a consumer: {}".format(op_name)))
585 586

    def push(self, channeldata, op_name=None):
B
barrierye 已提交
587
        _LOGGER.debug(
B
barriery 已提交
588 589
            self._log("(logid={}) Op({}) Pushing data".format(channeldata.id,
                                                              op_name)))
590
        if len(self._producers) == 0:
591
            _LOGGER.critical(
592
                self._log(
B
barriery 已提交
593 594 595
                    "(logid={}) Op({}) Failed to push data: expected number of "
                    "producers to be greater than 0, but the it is 0.".format(
                        channeldata.id, op_name)))
596
            os._exit(-1)
597 598 599 600
        elif len(self._producers) == 1:
            with self._cv:
                while self._stop is False:
                    try:
B
barrierye 已提交
601
                        self.put({op_name: channeldata}, timeout=0)
602 603 604
                        break
                    except Queue.Full:
                        self._cv.wait()
B
barrierye 已提交
605 606
                if self._stop:
                    raise ChannelStopError()
607
                self._cv.notify_all()
B
barriery 已提交
608
            _LOGGER.debug(
B
barriery 已提交
609 610
                self._log("(logid={}) Op({}) Pushed data into internal_queue.".
                          format(channeldata.id, op_name)))
611 612
            return True
        elif op_name is None:
613
            _LOGGER.critical(
614
                self._log(
B
barriery 已提交
615 616 617
                    "(logid={}) Op({}) Failed to push data: there are multiple"
                    " producers, so op_name cannot be None.".format(
                        channeldata.id, op_name)))
618
            os._exit(-1)
619 620 621 622 623

        producer_num = len(self._producers)
        data_id = channeldata.id
        put_data = None
        with self._cv:
B
barrierye 已提交
624 625
            if data_id not in self._input_buf:
                self._input_buf[data_id] = {
626 627 628
                    name: None
                    for name in self._producers
                }
B
barrierye 已提交
629
                self._pushed_producer_count[data_id] = 0
B
barrierye 已提交
630
            self._input_buf[data_id][op_name] = channeldata
B
barrierye 已提交
631
            if self._pushed_producer_count[data_id] + 1 == producer_num:
B
barrierye 已提交
632 633
                put_data = self._input_buf[data_id]
                self._input_buf.pop(data_id)
B
barrierye 已提交
634
                self._pushed_producer_count.pop(data_id)
635
            else:
B
barrierye 已提交
636
                self._pushed_producer_count[data_id] += 1
637 638

            if put_data is None:
B
barrierye 已提交
639
                _LOGGER.debug(
B
barriery 已提交
640 641 642
                    self._log(
                        "(logid={}) Op({}) Pushed data into input_buffer.".
                        format(data_id, op_name)))
643 644 645 646 647 648 649
            else:
                while self._stop is False:
                    try:
                        self.put(put_data, timeout=0)
                        break
                    except Queue.Empty:
                        self._cv.wait()
B
barrierye 已提交
650 651
                if self._stop:
                    raise ChannelStopError()
652

B
barrierye 已提交
653
                _LOGGER.debug(
B
barriery 已提交
654 655 656
                    self._log(
                        "(logid={}) Op({}) Pushed data into internal_queue.".
                        format(data_id, op_name)))
657 658 659
            self._cv.notify_all()
        return True

B
barriery 已提交
660
    def front(self, op_name=None, timeout=None):
B
barriery 已提交
661
        _LOGGER.debug(
B
barriery 已提交
662 663
            self._log("Op({}) Getting data[?]; timeout(s)={}".format(op_name,
                                                                     timeout)))
B
barriery 已提交
664
        endtime = None
B
bug fix  
barriery 已提交
665 666 667 668 669
        if timeout is not None:
            if timeout <= 0:
                timeout = None
            else:
                endtime = _time() + timeout
B
barriery 已提交
670

B
barrierye 已提交
671
        if len(self._consumer_cursors) == 0:
672
            _LOGGER.critical(
673
                self._log(
B
barriery 已提交
674 675
                    "Op({}) Failed to get data: expected number of consumers to be "
                    "greater than 0, but the it is 0.".format(op_name)))
676
            os._exit(-1)
B
barrierye 已提交
677
        elif len(self._consumer_cursors) == 1:
678 679 680 681 682 683 684
            resp = None
            with self._cv:
                while self._stop is False and resp is None:
                    try:
                        resp = self.get(timeout=0)
                        break
                    except Queue.Empty:
B
barriery 已提交
685 686 687
                        if timeout is not None:
                            remaining = endtime - _time()
                            if remaining <= 0.0:
B
barriery 已提交
688
                                _LOGGER.debug(
B
barriery 已提交
689 690 691
                                    self._log(
                                        "Op({}) Failed to get data: timeout".
                                        format(op_name)))
B
barriery 已提交
692 693 694 695
                                raise ChannelTimeoutError()
                            self._cv.wait(remaining)
                        else:
                            self._cv.wait()
B
barrierye 已提交
696 697
                if self._stop:
                    raise ChannelStopError()
B
barrierye 已提交
698
            _LOGGER.debug(
B
barriery 已提交
699 700
                self._log("(logid={}) Op({}) Got data".format(resp.values()[0]
                                                              .id, op_name)))
701 702
            return resp
        elif op_name is None:
703
            _LOGGER.critical(
B
barriery 已提交
704 705 706
                self._log("Op({}) Failed to get data: there are multiple "
                          "consumers, so op_name cannot be None.".format(
                              op_name)))
707
            os._exit(-1)
708

B
barrierye 已提交
709 710 711 712 713 714 715 716 717 718
        # In output_buf, different Ops (according to op_name) have different
        # cursors. In addition, there is a base_cursor. Their difference is
        # the data_idx to be taken by the corresponding Op at the current
        # time:    data_idx = consumer_cursor - base_cursor
        # 
        #            base_cursor    consumer_B_cursor (data_idx: 3)
        #                 |                       |
        # output_buf: | data0 | data1 | data2 | data3 |
        #                 |
        #   consumer_A_cursor (data_idx: 0)
719
        with self._cv:
B
barrierye 已提交
720 721 722 723
            # When the data required by the current Op is not in output_buf,
            # it is necessary to obtain a data from queue and add it to output_buf.
            while self._stop is False and self._consumer_cursors[
                    op_name] - self._base_cursor >= len(self._output_buf):
724 725
                try:
                    channeldata = self.get(timeout=0)
B
barrierye 已提交
726
                    self._output_buf.append(channeldata)
B
barriery 已提交
727
                    _LOGGER.debug(
B
barriery 已提交
728 729 730
                        self._log(
                            "(logid={}) Op({}) Pop ready item into output_buffer".
                            format(channeldata.values()[0].id, op_name)))
731 732
                    break
                except Queue.Empty:
B
barriery 已提交
733 734 735
                    if timeout is not None:
                        remaining = endtime - _time()
                        if remaining <= 0.0:
B
barriery 已提交
736
                            _LOGGER.debug(
B
barriery 已提交
737 738
                                self._log("Op({}) Failed to get data: timeout".
                                          format(op_name)))
B
barriery 已提交
739 740 741 742
                            raise ChannelTimeoutError()
                        self._cv.wait(remaining)
                    else:
                        self._cv.wait()
B
barrierye 已提交
743 744
            if self._stop:
                raise ChannelStopError()
745

B
barrierye 已提交
746 747 748
            consumer_cursor = self._consumer_cursors[op_name]
            base_cursor = self._base_cursor
            data_idx = consumer_cursor - base_cursor
B
barrierye 已提交
749 750

            resp = None
751

B
barrierye 已提交
752 753 754 755 756 757
            self._cursor_count[consumer_cursor] -= 1
            if consumer_cursor == base_cursor and self._cursor_count[
                    consumer_cursor] == 0:
                # When all the different Ops get the data that data_idx points
                # to, pop the data from output_buf.
                self._cursor_count.pop(consumer_cursor)
B
barrierye 已提交
758
                resp = self._output_buf.pop(0)
B
barrierye 已提交
759
                self._base_cursor += 1
760 761
                # to avoid cursor overflow
                if self._base_cursor >= self._reset_max_cursor:
B
barriery 已提交
762
                    _LOGGER.info(self._log("Reset cursor in Channel"))
763 764 765 766 767 768 769
                    self._base_cursor -= self._reset_max_cursor
                    for name in self._consumer_cursors:
                        self._consumer_cursors[name] -= self._reset_max_cursor
                    self._cursor_count = {
                        cursor - self._reset_max_cursor: count
                        for cursor, count in self._cursor_count.items()
                    }
B
barrierye 已提交
770 771
            else:
                resp = copy.deepcopy(self._output_buf[data_idx])
B
barrierye 已提交
772 773 774 775 776 777

            self._consumer_cursors[op_name] += 1
            new_consumer_cursor = self._consumer_cursors[op_name]
            if self._cursor_count.get(new_consumer_cursor) is None:
                self._cursor_count[new_consumer_cursor] = 0
            self._cursor_count[new_consumer_cursor] += 1
778 779 780

            self._cv.notify_all()

B
barriery 已提交
781
        _LOGGER.debug(
B
barriery 已提交
782 783
            self._log("(logid={}) Op({}) Got data from output_buffer".format(
                resp.values()[0].id, op_name)))
B
barrierye 已提交
784
        return resp
785 786

    def stop(self):
787
        _LOGGER.info(self._log("stop."))
788
        self._stop = True
B
barrierye 已提交
789 790 791
        with self._cv:
            self._cv.notify_all()

B
barriery 已提交
792

B
barriery 已提交
793 794 795
class ChannelTimeoutError(RuntimeError):
    def __init__(self):
        pass
B
barrierye 已提交
796

B
barriery 已提交
797

B
barrierye 已提交
798 799 800
class ChannelStopError(RuntimeError):
    def __init__(self):
        pass