channel.py 34.0 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

T
TeslaZhao 已提交
35 36 37 38
class ChannelDataErrcode(enum.Enum):
    """
    ChannelData error code
    """
D
dongdaxiang 已提交
39 40 41 42 43
    OK = 0
    TIMEOUT = 1
    NOT_IMPLEMENTED = 2
    TYPE_ERROR = 3
    RPC_PACKAGE_ERROR = 4
B
barrierye 已提交
44
    CLIENT_ERROR = 5
B
barrierye 已提交
45
    CLOSED_ERROR = 6
B
barriery 已提交
46 47
    NO_SERVICE = 7
    UNKNOW = 8
T
TeslaZhao 已提交
48 49 50 51 52 53 54 55 56
    PRODUCT_ERROR = 9


class ProductErrCode(enum.Enum):
    """
    ProductErrCode is a base class for recording business error code. 
    product developers inherit this class and extend more error codes. 
    """
    pass
D
dongdaxiang 已提交
57 58 59


class ChannelDataType(enum.Enum):
60 61 62
    """
    Channel data type
    """
D
dongdaxiang 已提交
63 64 65 66 67
    DICT = 0
    CHANNEL_NPDATA = 1
    ERROR = 2


68 69 70 71
class ChannelData(object):
    def __init__(self,
                 datatype=None,
                 npdata=None,
B
barrierye 已提交
72
                 dictdata=None,
73
                 data_id=None,
T
TeslaZhao 已提交
74 75
                 log_id=None,
                 error_code=None,
B
barrierye 已提交
76
                 error_info=None,
T
TeslaZhao 已提交
77 78
                 prod_error_code=None,
                 prod_error_info=None,
B
barrierye 已提交
79
                 client_need_profile=False):
80 81 82
        '''
        There are several ways to use it:
        
T
TeslaZhao 已提交
83 84 85
        1. ChannelData(ChannelDataType.CHANNEL_NPDATA.value, npdata, data_id, log_id)
        2. ChannelData(ChannelDataType.DICT.value, dictdata, data_id, log_id)
        3. ChannelData(error_code, error_info, prod_error_code, prod_error_info, data_id, log_id)
86 87 88 89

        Protobufs are not pickle-able:
        https://stackoverflow.com/questions/55344376/how-to-import-protobuf-module
        '''
T
TeslaZhao 已提交
90
        if error_code is not None or prod_error_code is not None:
91
            if data_id is None or error_info is None:
B
barriery 已提交
92 93
                _LOGGER.critical("Failed to generate ChannelData: data_id"
                                 " and error_info cannot be None")
94
                os._exit(-1)
95 96
            datatype = ChannelDataType.ERROR.value
        else:
B
barrierye 已提交
97
            if datatype == ChannelDataType.CHANNEL_NPDATA.value:
T
TeslaZhao 已提交
98 99
                error_code, error_info = ChannelData.check_npdata(npdata)
                if error_code != ChannelDataErrcode.OK.value:
B
barrierye 已提交
100
                    datatype = ChannelDataType.ERROR.value
T
TeslaZhao 已提交
101 102
                    _LOGGER.error("(data_id={} log_id={}) {}".format(
                        data_id, log_id, error_info))
B
barrierye 已提交
103
            elif datatype == ChannelDataType.DICT.value:
T
TeslaZhao 已提交
104 105
                error_code, error_info = ChannelData.check_dictdata(dictdata)
                if error_code != ChannelDataErrcode.OK.value:
B
barrierye 已提交
106
                    datatype = ChannelDataType.ERROR.value
T
TeslaZhao 已提交
107 108
                    _LOGGER.error("(data_id={} log_id={}) {}".format(
                        data_id, log_id, error_info))
109
            else:
T
TeslaZhao 已提交
110 111
                _LOGGER.critical("(data_id={} log_id={}) datatype not match".
                                 format(data_id, log_id))
112
                os._exit(-1)
113
        self.datatype = datatype
B
barrierye 已提交
114 115
        self.npdata = npdata
        self.dictdata = dictdata
116
        self.id = data_id
T
TeslaZhao 已提交
117 118
        self.log_id = log_id
        self.error_code = error_code
119
        self.error_info = error_info
T
TeslaZhao 已提交
120 121
        self.prod_error_code = prod_error_code
        self.prod_error_info = prod_error_info
B
barrierye 已提交
122
        self.client_need_profile = client_need_profile
B
barrierye 已提交
123
        self.profile_data_set = set()
B
barrierye 已提交
124

B
barrierye 已提交
125
    def add_profile(self, profile_set):
B
barrierye 已提交
126 127
        if self.client_need_profile is False:
            self.client_need_profile = True
B
barrierye 已提交
128
        self.profile_data_set |= profile_set
129

B
barrierye 已提交
130 131
    @staticmethod
    def check_dictdata(dictdata):
T
TeslaZhao 已提交
132
        error_code = ChannelDataErrcode.OK.value
B
barrierye 已提交
133
        error_info = None
B
barrierye 已提交
134 135 136 137
        if isinstance(dictdata, list):
            # batch data
            for sample in dictdata:
                if not isinstance(sample, dict):
T
TeslaZhao 已提交
138
                    error_code = ChannelDataErrcode.TYPE_ERROR.value
B
barriery 已提交
139 140
                    error_info = "Failed to check data: the type of " \
                            "data must be dict, but get {}.".format(type(sample))
B
barrierye 已提交
141 142 143
                    break
        elif not isinstance(dictdata, dict):
            # batch size = 1
T
TeslaZhao 已提交
144
            error_code = ChannelDataErrcode.TYPE_ERROR.value
B
barriery 已提交
145 146
            error_info = "Failed to check data: the type of data must " \
                    "be dict, but get {}.".format(type(dictdata))
T
TeslaZhao 已提交
147
        return error_code, error_info
B
barrierye 已提交
148

B
bug fix  
barriery 已提交
149 150
    @staticmethod
    def check_batch_npdata(batch):
T
TeslaZhao 已提交
151
        error_code = ChannelDataErrcode.OK.value
B
bug fix  
barriery 已提交
152 153
        error_info = None
        for npdata in batch:
T
TeslaZhao 已提交
154 155
            error_code, error_info = ChannelData.check_npdata(npdata)
            if error_code != ChannelDataErrcode.OK.value:
B
bug fix  
barriery 已提交
156
                break
T
TeslaZhao 已提交
157
        return error_code, error_info
B
bug fix  
barriery 已提交
158

B
barrierye 已提交
159 160
    @staticmethod
    def check_npdata(npdata):
T
TeslaZhao 已提交
161
        error_code = ChannelDataErrcode.OK.value
162
        error_info = None
W
wangjiawei04 已提交
163 164 165 166
        if isinstance(npdata, list):
            # batch data
            for sample in npdata:
                if not isinstance(sample, dict):
T
TeslaZhao 已提交
167
                    error_code = ChannelDataErrcode.TYPE_ERROR.value
B
barriery 已提交
168 169 170
                    error_info = "Failed to check data: the " \
                            "value of data must be dict, but get {}.".format(
                                    type(sample))
W
wangjiawei04 已提交
171 172 173
                    break
                for _, value in sample.items():
                    if not isinstance(value, np.ndarray):
T
TeslaZhao 已提交
174
                        error_code = ChannelDataErrcode.TYPE_ERROR.value
B
barriery 已提交
175 176 177
                        error_info = "Failed to check data: the" \
                                " value of data must be np.ndarray, but get {}.".format(
                                        type(value))
T
TeslaZhao 已提交
178
                        return error_code, error_info
W
wangjiawei04 已提交
179 180 181 182
        elif isinstance(npdata, dict):
            # batch_size = 1
            for _, value in npdata.items():
                if not isinstance(value, np.ndarray):
T
TeslaZhao 已提交
183
                    error_code = ChannelDataErrcode.TYPE_ERROR.value
B
barriery 已提交
184 185 186
                    error_info = "Failed to check data: the value " \
                            "of data must be np.ndarray, but get {}.".format(
                                    type(value))
W
wangjiawei04 已提交
187 188
                    break
        else:
T
TeslaZhao 已提交
189
            error_code = ChannelDataErrcode.TYPE_ERROR.value
B
barriery 已提交
190 191
            error_info = "Failed to check data: the value of data " \
                    "must be dict, but get {}.".format(type(npdata))
T
TeslaZhao 已提交
192
        return error_code, error_info
193 194 195

    def parse(self):
        feed = None
B
barrierye 已提交
196 197
        if self.datatype == ChannelDataType.CHANNEL_NPDATA.value:
            # return narray
198
            feed = self.npdata
B
barrierye 已提交
199 200 201
        elif self.datatype == ChannelDataType.DICT.value:
            # return dict
            feed = self.dictdata
202
        else:
B
barriery 已提交
203 204
            _LOGGER.critical("Failed to parse channeldata: error " \
                    "type({}) in datatype.".format(self.datatype))
205
            os._exit(-1)
206 207
        return feed

208 209 210 211 212 213 214 215
    def __cmp__(self, other):
        if self.id < other.id:
            return -1
        elif self.id == other.id:
            return 0
        else:
            return 1

216
    def __str__(self):
217
        return "type[{}], error_code[{}], data_id[{}], log_id[{}], dict_data[{}]".format(
T
TeslaZhao 已提交
218
            ChannelDataType(self.datatype).name, self.error_code, self.id,
219
            self.log_id, str(self.dictdata))
220 221


B
barrierye 已提交
222
class ProcessChannel(object):
223 224 225 226 227 228 229 230 231
    """ 
    (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 已提交
232
    3. Function front support timeout param to make auto-batching.
233 234 235 236 237

    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 已提交
238 239 240 241 242 243 244 245 246 247 248

    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.
249 250
    """

B
barriery 已提交
251
    def __init__(self, manager, name=None, maxsize=0):
B
barrierye 已提交
252 253 254 255 256 257
        # 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
258
        self._que = manager.PriorityQueue(maxsize=maxsize)
259 260
        self._maxsize = maxsize
        self.name = name
261
        self._stop = manager.Value('i', 0)
262 263 264 265

        self._cv = multiprocessing.Condition()

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

269
        self._reset_max_cursor = 1000000000000000000
B
barrierye 已提交
270 271 272 273
        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()
274

B
barriery 已提交
275 276 277
    def get_maxsize(self):
        return self._maxsize

B
barriery 已提交
278 279 280
    def size(self):
        return self._que.qsize()

281 282 283 284
    def get_producers(self):
        return self._producers

    def get_consumers(self):
B
barrierye 已提交
285
        return self._consumer_cursors.keys()
286 287 288 289 290 291 292

    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:
293
            _LOGGER.critical(
B
barriery 已提交
294 295
                self._log("Failed to add producer: producer({})" \
                        " is already in channel".format(op_name)))
296
            os._exit(-1)
297
        self._producers.append(op_name)
B
barriery 已提交
298
        _LOGGER.debug(self._log("Succ add a producer: {}".format(op_name)))
299 300 301

    def add_consumer(self, op_name):
        """ not thread safe, and can only be called during initialization. """
B
barrierye 已提交
302
        if op_name in self._consumer_cursors:
303
            _LOGGER.critical(
B
barriery 已提交
304 305
                    self._log("Failed to add consumer: consumer({})" \
                            " is already in channel".format(op_name)))
306
            os._exit(-1)
B
barrierye 已提交
307
        self._consumer_cursors[op_name] = 0
308

B
barrierye 已提交
309 310 311
        if self._cursor_count.get(0) is None:
            self._cursor_count[0] = 0
        self._cursor_count[0] += 1
B
barriery 已提交
312
        _LOGGER.debug(self._log("Succ add a consumer: {}".format(op_name)))
313 314

    def push(self, channeldata, op_name=None):
B
barrierye 已提交
315
        _LOGGER.debug(
316 317
            self._log("(data_id={} log_id={}) Op({}) Enter channel::push".
                      format(channeldata.id, channeldata.log_id, op_name)))
318
        if len(self._producers) == 0:
319
            _LOGGER.critical(
320
                self._log(
T
TeslaZhao 已提交
321
                    "(data_id={} log_id={}) Op({}) Failed to push data: expected number"
B
barriery 已提交
322
                    " of producers to be greater than 0, but the it is 0.".
T
TeslaZhao 已提交
323
                    format(channeldata.id, channeldata.log_id, op_name)))
324
            os._exit(-1)
325 326
        elif len(self._producers) == 1:
            with self._cv:
327
                while self._stop.value == 0:
328
                    try:
T
TeslaZhao 已提交
329 330 331 332
                        self._que.put((channeldata.id, {
                            op_name: channeldata
                        }),
                                      timeout=0)
333 334 335
                        break
                    except Queue.Full:
                        self._cv.wait()
336
                if self._stop.value == 1:
B
barrierye 已提交
337
                    raise ChannelStopError()
338
                self._cv.notify_all()
B
barriery 已提交
339
            _LOGGER.debug(
T
TeslaZhao 已提交
340 341 342
                self._log(
                    "(data_id={} log_id={}) Op({}) Pushed data into internal queue.".
                    format(channeldata.id, channeldata.log_id, op_name)))
343 344
            return True
        elif op_name is None:
345
            _LOGGER.critical(
346
                self._log(
T
TeslaZhao 已提交
347
                    "(data_id={} log_id={}) Op({}) Failed to push data: there are multiple "
B
barriery 已提交
348
                    "producers, so op_name cannot be None.".format(
T
TeslaZhao 已提交
349
                        channeldata.id, channeldata.log_id, op_name)))
350
            os._exit(-1)
351 352 353

        producer_num = len(self._producers)
        data_id = channeldata.id
T
TeslaZhao 已提交
354
        log_id = channeldata.log_id
355 356
        put_data = None
        with self._cv:
B
barrierye 已提交
357 358
            if data_id not in self._input_buf:
                self._input_buf[data_id] = {
359 360 361
                    name: None
                    for name in self._producers
                }
B
barrierye 已提交
362
                self._pushed_producer_count[data_id] = 0
363
            # see: https://docs.python.org/3.6/library/multiprocessing.html?highlight=multiprocess#proxy-objects
B
barrierye 已提交
364 365 366 367 368
            # 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 已提交
369
            if self._pushed_producer_count[data_id] + 1 == producer_num:
B
barrierye 已提交
370 371
                put_data = self._input_buf[data_id]
                self._input_buf.pop(data_id)
B
barrierye 已提交
372
                self._pushed_producer_count.pop(data_id)
373
            else:
B
barrierye 已提交
374
                self._pushed_producer_count[data_id] += 1
375 376

            if put_data is None:
B
barrierye 已提交
377
                _LOGGER.debug(
B
barriery 已提交
378
                    self._log(
T
TeslaZhao 已提交
379 380
                        "(data_id={} log_id={}) Op({}) Pushed data into input_buffer.".
                        format(data_id, log_id, op_name)))
381
            else:
382
                while self._stop.value == 0:
383
                    try:
T
TeslaZhao 已提交
384
                        self._que.put((data_id, put_data), timeout=0)
385 386 387
                        break
                    except Queue.Empty:
                        self._cv.wait()
388
                if self._stop.value == 1:
B
barrierye 已提交
389
                    raise ChannelStopError()
390

B
barrierye 已提交
391
                _LOGGER.debug(
B
barriery 已提交
392
                    self._log(
T
TeslaZhao 已提交
393 394
                        "(data_id={} log_id={}) Op({}) Pushed data into internal_queue.".
                        format(data_id, log_id, op_name)))
395 396 397
            self._cv.notify_all()
        return True

B
barriery 已提交
398
    def front(self, op_name=None, timeout=None):
B
barriery 已提交
399
        _LOGGER.debug(
B
barriery 已提交
400 401
            self._log("Op({}) Getting data[?]; timeout(s)={}".format(op_name,
                                                                     timeout)))
B
barriery 已提交
402
        endtime = None
B
bug fix  
barriery 已提交
403 404 405 406 407
        if timeout is not None:
            if timeout <= 0:
                timeout = None
            else:
                endtime = _time() + timeout
B
barriery 已提交
408

B
barrierye 已提交
409
        if len(self._consumer_cursors) == 0:
410
            _LOGGER.critical(
411
                self._log(
B
barriery 已提交
412 413
                    "Op({}) Failed to get data: expected number of consumers to be " \
                            "greater than 0, but the it is 0.".format(op_name)))
414
            os._exit(-1)
B
barrierye 已提交
415
        elif len(self._consumer_cursors) == 1:
416 417
            resp = None
            with self._cv:
418
                while self._stop.value == 0 and resp is None:
419
                    try:
T
TeslaZhao 已提交
420
                        resp = self._que.get(timeout=0)[1]
421 422
                        break
                    except Queue.Empty:
B
barriery 已提交
423 424 425
                        if timeout is not None:
                            remaining = endtime - _time()
                            if remaining <= 0.0:
B
barriery 已提交
426
                                _LOGGER.debug(
B
barriery 已提交
427 428
                                    self._log("Op({}) Failed to get data: "
                                              "timeout".format(op_name)))
B
barriery 已提交
429 430 431 432
                                raise ChannelTimeoutError()
                            self._cv.wait(remaining)
                        else:
                            self._cv.wait()
433
                if self._stop.value == 1:
B
barrierye 已提交
434
                    raise ChannelStopError()
T
TeslaZhao 已提交
435 436 437 438 439 440

            if resp is not None:
                list_values = list(resp.values())
                _LOGGER.debug(
                    self._log("(data_id={} log_id={}) Op({}) Got data".format(
                        list_values[0].id, list_values[0].log_id, op_name)))
441 442
            return resp
        elif op_name is None:
443
            _LOGGER.critical(
444
                self._log(
B
barriery 已提交
445 446
                    "Op({}) Failed to get data: there are multiple consumers, "
                    "so op_name cannot be None.".format(op_name)))
447
            os._exit(-1)
448

B
barrierye 已提交
449 450 451 452 453 454 455 456 457 458
        # 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)
459
        with self._cv:
B
barrierye 已提交
460 461
            # 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.
462
            while self._stop.value == 0 and self._consumer_cursors[
B
barrierye 已提交
463
                    op_name] - self._base_cursor.value >= len(self._output_buf):
464
                try:
T
TeslaZhao 已提交
465
                    channeldata = self._que.get(timeout=0)[1]
B
barrierye 已提交
466
                    self._output_buf.append(channeldata)
T
TeslaZhao 已提交
467
                    list_values = list(channeldata.values())
B
barriery 已提交
468
                    _LOGGER.debug(
B
barriery 已提交
469
                        self._log(
T
TeslaZhao 已提交
470
                            "(data_id={} log_id={}) Op({}) Pop ready item into output_buffer".
T
TeslaZhao 已提交
471 472
                            format(list_values[0].id, list_values[0].log_id,
                                   op_name)))
473 474
                    break
                except Queue.Empty:
B
barriery 已提交
475 476 477
                    if timeout is not None:
                        remaining = endtime - _time()
                        if remaining <= 0.0:
B
barriery 已提交
478
                            _LOGGER.debug(
B
barriery 已提交
479 480
                                self._log("Op({}) Failed to get data: timeout".
                                          format(op_name)))
B
barriery 已提交
481 482 483 484
                            raise ChannelTimeoutError()
                        self._cv.wait(remaining)
                    else:
                        self._cv.wait()
485
            if self._stop.value == 1:
B
barrierye 已提交
486
                raise ChannelStopError()
487

B
barrierye 已提交
488 489 490 491
            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]
492

B
barrierye 已提交
493 494 495 496 497 498 499 500
            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
501 502
                # to avoid cursor overflow
                if self._base_cursor.value >= self._reset_max_cursor:
B
barriery 已提交
503
                    _LOGGER.info(self._log("Reset cursor in Channel"))
504 505 506 507 508 509 510 511 512 513
                    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 已提交
514 515 516 517 518 519

            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
520

521 522
            self._cv.notify_all()

T
TeslaZhao 已提交
523 524 525 526 527 528
        if resp is not None:
            list_values = list(resp.values())
            _LOGGER.debug(
                self._log(
                    "(data_id={} log_id={}) Op({}) Got data from output_buffer".
                    format(list_values[0].id, list_values[0].log_id, op_name)))
B
barriery 已提交
529
        return resp
530 531

    def stop(self):
532
        _LOGGER.info(self._log("stop."))
533
        self._stop.value = 1
B
barrierye 已提交
534 535
        with self._cv:
            self._cv.notify_all()
536 537


538
class ThreadChannel(Queue.PriorityQueue):
539 540 541 542 543 544 545 546 547
    """ 
    (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 已提交
548
    3. Function front support timeout param to make auto-batching.
549 550 551 552 553

    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 已提交
554 555 556 557 558 559 560 561 562 563 564

    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.
565 566
    """

B
barriery 已提交
567
    def __init__(self, name=None, maxsize=-1):
568 569 570 571 572 573 574 575
        Queue.Queue.__init__(self, maxsize=maxsize)
        self._maxsize = maxsize
        self.name = name
        self._stop = False

        self._cv = threading.Condition()

        self._producers = []
B
barrierye 已提交
576
        self._pushed_producer_count = {}  # {data_id: count}
B
barrierye 已提交
577
        self._input_buf = {}  # {data_id: {op_name: data}}
578

579
        self._reset_max_cursor = 1000000000000000000
B
barrierye 已提交
580 581 582 583
        self._consumer_cursors = {}  # {op_name: idx}
        self._cursor_count = {}  # {cursor: count}
        self._base_cursor = 0
        self._output_buf = []
584

B
barriery 已提交
585 586 587
    def get_maxsize(self):
        return self._maxsize

B
barriery 已提交
588 589 590
    def size(self):
        return self.qsize()

591 592 593 594
    def get_producers(self):
        return self._producers

    def get_consumers(self):
B
barrierye 已提交
595
        return self._consumer_cursors.keys()
596 597 598 599 600 601 602

    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:
603
            _LOGGER.critical(
B
barriery 已提交
604 605
                self._log("Failed to add producer: producer({}) is "
                          "already in channel".format(op_name)))
606
            os._exit(-1)
607
        self._producers.append(op_name)
B
barriery 已提交
608
        _LOGGER.debug(self._log("Succ add a producer: {}".format(op_name)))
609 610 611

    def add_consumer(self, op_name):
        """ not thread safe, and can only be called during initialization. """
B
barrierye 已提交
612
        if op_name in self._consumer_cursors:
613
            _LOGGER.critical(
B
barriery 已提交
614 615
                self._log("Failed to add consumer: consumer({}) is "
                          "already in channel".format(op_name)))
616
            os._exit(-1)
B
barrierye 已提交
617
        self._consumer_cursors[op_name] = 0
618

B
barrierye 已提交
619 620 621
        if self._cursor_count.get(0) is None:
            self._cursor_count[0] = 0
        self._cursor_count[0] += 1
B
barriery 已提交
622
        _LOGGER.debug(self._log("Succ add a consumer: {}".format(op_name)))
623 624

    def push(self, channeldata, op_name=None):
B
barrierye 已提交
625
        _LOGGER.debug(
T
TeslaZhao 已提交
626 627
            self._log("(data_id={} log_id={}) Op({}) Pushing data".format(
                channeldata.id, channeldata.log_id, op_name)))
628
        if len(self._producers) == 0:
629
            _LOGGER.critical(
630
                self._log(
T
TeslaZhao 已提交
631
                    "(data_id={} log_id={}) Op({}) Failed to push data: expected number of "
B
barriery 已提交
632
                    "producers to be greater than 0, but the it is 0.".format(
T
TeslaZhao 已提交
633
                        channeldata.id, channeldata.log_id, op_name)))
634
            os._exit(-1)
635 636 637 638
        elif len(self._producers) == 1:
            with self._cv:
                while self._stop is False:
                    try:
T
TeslaZhao 已提交
639 640 641 642
                        self.put((channeldata.id, {
                            op_name: channeldata
                        }),
                                 timeout=0)
643 644 645
                        break
                    except Queue.Full:
                        self._cv.wait()
B
barrierye 已提交
646 647
                if self._stop:
                    raise ChannelStopError()
648
                self._cv.notify_all()
B
barriery 已提交
649
            _LOGGER.debug(
T
TeslaZhao 已提交
650 651 652
                self._log(
                    "(data_id={} log_id={}) Op({}) Pushed data into internal_queue.".
                    format(channeldata.id, channeldata.log_id, op_name)))
653 654
            return True
        elif op_name is None:
655
            _LOGGER.critical(
656
                self._log(
T
TeslaZhao 已提交
657
                    "(data_id={} log_id={}) Op({}) Failed to push data: there are multiple"
B
barriery 已提交
658
                    " producers, so op_name cannot be None.".format(
T
TeslaZhao 已提交
659
                        channeldata.id, channeldata.log_id, op_name)))
660
            os._exit(-1)
661 662 663

        producer_num = len(self._producers)
        data_id = channeldata.id
T
TeslaZhao 已提交
664
        log_id = channeldata.log_id
665 666
        put_data = None
        with self._cv:
B
barrierye 已提交
667 668
            if data_id not in self._input_buf:
                self._input_buf[data_id] = {
669 670 671
                    name: None
                    for name in self._producers
                }
B
barrierye 已提交
672
                self._pushed_producer_count[data_id] = 0
B
barrierye 已提交
673
            self._input_buf[data_id][op_name] = channeldata
B
barrierye 已提交
674
            if self._pushed_producer_count[data_id] + 1 == producer_num:
B
barrierye 已提交
675 676
                put_data = self._input_buf[data_id]
                self._input_buf.pop(data_id)
B
barrierye 已提交
677
                self._pushed_producer_count.pop(data_id)
678
            else:
B
barrierye 已提交
679
                self._pushed_producer_count[data_id] += 1
680 681

            if put_data is None:
B
barrierye 已提交
682
                _LOGGER.debug(
B
barriery 已提交
683
                    self._log(
T
TeslaZhao 已提交
684 685
                        "(data_id={} log_id={}) Op({}) Pushed data into input_buffer.".
                        format(data_id, log_id, op_name)))
686 687 688
            else:
                while self._stop is False:
                    try:
T
TeslaZhao 已提交
689
                        self.put((data_id, put_data), timeout=0)
690 691 692
                        break
                    except Queue.Empty:
                        self._cv.wait()
B
barrierye 已提交
693 694
                if self._stop:
                    raise ChannelStopError()
695

B
barrierye 已提交
696
                _LOGGER.debug(
B
barriery 已提交
697
                    self._log(
T
TeslaZhao 已提交
698 699
                        "(data_id={} log_id={}) Op({}) Pushed data into internal_queue.".
                        format(data_id, log_id, op_name)))
700 701 702
            self._cv.notify_all()
        return True

B
barriery 已提交
703
    def front(self, op_name=None, timeout=None):
B
barriery 已提交
704
        _LOGGER.debug(
B
barriery 已提交
705 706
            self._log("Op({}) Getting data[?]; timeout(s)={}".format(op_name,
                                                                     timeout)))
B
barriery 已提交
707
        endtime = None
B
bug fix  
barriery 已提交
708 709 710 711 712
        if timeout is not None:
            if timeout <= 0:
                timeout = None
            else:
                endtime = _time() + timeout
B
barriery 已提交
713

B
barrierye 已提交
714
        if len(self._consumer_cursors) == 0:
715
            _LOGGER.critical(
716
                self._log(
B
barriery 已提交
717 718
                    "Op({}) Failed to get data: expected number of consumers to be "
                    "greater than 0, but the it is 0.".format(op_name)))
719
            os._exit(-1)
B
barrierye 已提交
720
        elif len(self._consumer_cursors) == 1:
721 722 723 724
            resp = None
            with self._cv:
                while self._stop is False and resp is None:
                    try:
T
TeslaZhao 已提交
725
                        resp = self.get(timeout=0)[1]
726 727
                        break
                    except Queue.Empty:
B
barriery 已提交
728 729 730
                        if timeout is not None:
                            remaining = endtime - _time()
                            if remaining <= 0.0:
B
barriery 已提交
731
                                _LOGGER.debug(
B
barriery 已提交
732 733 734
                                    self._log(
                                        "Op({}) Failed to get data: timeout".
                                        format(op_name)))
B
barriery 已提交
735 736 737 738
                                raise ChannelTimeoutError()
                            self._cv.wait(remaining)
                        else:
                            self._cv.wait()
B
barrierye 已提交
739 740
                if self._stop:
                    raise ChannelStopError()
T
TeslaZhao 已提交
741 742 743 744 745
            if resp is not None:
                list_values = list(resp.values())
                _LOGGER.debug(
                    self._log("(data_id={} log_id={}) Op({}) Got data".format(
                        list_values[0].id, list_values[0].log_id, op_name)))
746 747
            return resp
        elif op_name is None:
748
            _LOGGER.critical(
B
barriery 已提交
749 750 751
                self._log("Op({}) Failed to get data: there are multiple "
                          "consumers, so op_name cannot be None.".format(
                              op_name)))
752
            os._exit(-1)
753

B
barrierye 已提交
754 755 756 757 758 759 760 761 762 763
        # 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)
764
        with self._cv:
B
barrierye 已提交
765 766 767 768
            # 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):
769 770
                try:
                    channeldata = self.get(timeout=0)
B
barrierye 已提交
771
                    self._output_buf.append(channeldata)
T
TeslaZhao 已提交
772
                    list_values = list(channeldata.values())
B
barriery 已提交
773
                    _LOGGER.debug(
B
barriery 已提交
774
                        self._log(
T
TeslaZhao 已提交
775
                            "(data_id={} log_id={}) Op({}) Pop ready item into output_buffer".
T
TeslaZhao 已提交
776 777
                            format(list_values[0].id, list_values[0].log_id,
                                   op_name)))
778 779
                    break
                except Queue.Empty:
B
barriery 已提交
780 781 782
                    if timeout is not None:
                        remaining = endtime - _time()
                        if remaining <= 0.0:
B
barriery 已提交
783
                            _LOGGER.debug(
B
barriery 已提交
784 785
                                self._log("Op({}) Failed to get data: timeout".
                                          format(op_name)))
B
barriery 已提交
786 787 788 789
                            raise ChannelTimeoutError()
                        self._cv.wait(remaining)
                    else:
                        self._cv.wait()
B
barrierye 已提交
790 791
            if self._stop:
                raise ChannelStopError()
792

B
barrierye 已提交
793 794 795
            consumer_cursor = self._consumer_cursors[op_name]
            base_cursor = self._base_cursor
            data_idx = consumer_cursor - base_cursor
B
barrierye 已提交
796 797

            resp = None
798

B
barrierye 已提交
799 800 801 802 803 804
            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 已提交
805
                resp = self._output_buf.pop(0)
B
barrierye 已提交
806
                self._base_cursor += 1
807 808
                # to avoid cursor overflow
                if self._base_cursor >= self._reset_max_cursor:
B
barriery 已提交
809
                    _LOGGER.info(self._log("Reset cursor in Channel"))
810 811 812 813 814 815 816
                    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 已提交
817 818
            else:
                resp = copy.deepcopy(self._output_buf[data_idx])
B
barrierye 已提交
819 820 821 822 823 824

            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
825 826 827

            self._cv.notify_all()

T
TeslaZhao 已提交
828 829 830 831 832 833
        if resp is not None:
            list_values = list(resp.values())
            _LOGGER.debug(
                self._log(
                    "(data_id={} log_id={}) Op({}) Got data from output_buffer".
                    format(list_values[0].id, list_values[0].log_id, op_name)))
B
barrierye 已提交
834
        return resp
835 836

    def stop(self):
837
        _LOGGER.info(self._log("stop."))
838
        self._stop = True
B
barrierye 已提交
839 840 841
        with self._cv:
            self._cv.notify_all()

B
barriery 已提交
842

B
barriery 已提交
843 844 845
class ChannelTimeoutError(RuntimeError):
    def __init__(self):
        pass
B
barrierye 已提交
846

B
barriery 已提交
847

B
barrierye 已提交
848 849 850
class ChannelStopError(RuntimeError):
    def __init__(self):
        pass