oasislmf.pytools.common.event_stream ==================================== .. py:module:: oasislmf.pytools.common.event_stream .. autoapi-nested-parse:: Contain all common function and attribute to help read the event stream containing the losses Attributes ---------- .. autoapisummary:: oasislmf.pytools.common.event_stream.PIPE_CAPACITY oasislmf.pytools.common.event_stream.CDF_STREAM_ID oasislmf.pytools.common.event_stream.GUL_STREAM_ID oasislmf.pytools.common.event_stream.FM_STREAM_ID oasislmf.pytools.common.event_stream.LOSS_STREAM_ID oasislmf.pytools.common.event_stream.SUMMARY_STREAM_ID oasislmf.pytools.common.event_stream.ITEM_STREAM oasislmf.pytools.common.event_stream.COVERAGE_STREAM oasislmf.pytools.common.event_stream.NUM_SPECIAL_SIDX oasislmf.pytools.common.event_stream.MEAN_IDX oasislmf.pytools.common.event_stream.STD_DEV_IDX oasislmf.pytools.common.event_stream.TIV_IDX oasislmf.pytools.common.event_stream.MAX_LOSS_IDX Classes ------- .. autoapisummary:: oasislmf.pytools.common.event_stream.EventReader Functions --------- .. autoapisummary:: oasislmf.pytools.common.event_stream.stream_info_to_bytes oasislmf.pytools.common.event_stream.bytes_to_stream_types oasislmf.pytools.common.event_stream.read_stream_info oasislmf.pytools.common.event_stream.get_streams_in oasislmf.pytools.common.event_stream.get_and_check_header_in oasislmf.pytools.common.event_stream.init_streams_in oasislmf.pytools.common.event_stream.mv_read oasislmf.pytools.common.event_stream.reservation_overflows oasislmf.pytools.common.event_stream.mv_write oasislmf.pytools.common.event_stream.mv_write_summary_header oasislmf.pytools.common.event_stream.mv_write_summary_header_cached oasislmf.pytools.common.event_stream.mv_write_item_header oasislmf.pytools.common.event_stream.mv_write_sidx_loss oasislmf.pytools.common.event_stream.mv_write_sidx_loss_cached oasislmf.pytools.common.event_stream.write_mv_to_stream Module Contents --------------- .. py:data:: PIPE_CAPACITY :value: 65536 .. py:data:: CDF_STREAM_ID :value: 0 .. py:data:: GUL_STREAM_ID :value: 1 .. py:data:: FM_STREAM_ID :value: 2 .. py:data:: LOSS_STREAM_ID :value: 2 .. py:data:: SUMMARY_STREAM_ID :value: 3 .. py:data:: ITEM_STREAM :value: 1 .. py:data:: COVERAGE_STREAM :value: 2 .. py:data:: NUM_SPECIAL_SIDX :value: 5 .. py:data:: MEAN_IDX :value: -1 .. py:data:: STD_DEV_IDX :value: -2 .. py:data:: TIV_IDX :value: -3 .. py:data:: MAX_LOSS_IDX :value: -5 .. py:function:: stream_info_to_bytes(stream_source_type, stream_agg_type) From Stream source type and aggregation type produce the stream header :param stream_source_type: id of the tool that produced the stream :type stream_source_type: np.int32 :param stream_agg_type: id of the aggregation level of the stream :type stream_agg_type: np.int32 :returns: return bytes .. py:function:: bytes_to_stream_types(stream_header) Read the stream header and return the information on stream type :param stream_header: bytes :returns: (stream source type (np.int32), stream aggregation type (np.int32)) .. py:function:: read_stream_info(stream_obj) From open stream object return the information that characterize the stream (stream_source_type, stream_agg_type, len_sample) :param stream_obj: open stream :returns: (stream_source_type, stream_agg_type, len_sample) as np.int32 triplet .. py:function:: get_streams_in(files_in, stack) .. py:function:: get_and_check_header_in(streams_in) .. py:function:: init_streams_in(files_in, stack) If files_in use stdin as stream in otherwise open each path in files_in, read the header, check that they are the same, and return the streams and their info :param files_in: none or a list of path :param stack: contextlib stack to add the open stream to :returns: list of open streams and their info .. py:function:: mv_read(byte_mv, cursor, _dtype, itemsize) Read a certain dtype from numpy byte view starting at cursor, return the value and the index of the end of the object :param byte_mv: numpy byte view :param cursor: index of where the object start :param _dtype: data type of the object :param itemsize: size of the data type :returns: (object value, end of object index) .. py:function:: reservation_overflows(idx, reservation, capacity, name) Check whether reserving `reservation` more rows in a fixed-size output buffer (used before reading a summary's raw data, so its worst-case output is guaranteed to fit - see elt/manager.py and plt/manager.py's read_buffer) would overflow it. :param idx: current write position in the output buffer :param reservation: worst-case number of rows this summary could add :param capacity: output buffer's total size :param name: output type name, used only in the raised error message (e.g. "SELT") :returns: True if the buffer needs to be flushed before this summary can be read. False if there's room. :rtype: bool :raises ValueError: if idx == 0 (buffer already empty) and it still doesn't fit - flushing an empty buffer can never make more room, so returning True here would have the caller loop forever instead of making progress. .. py:function:: mv_write(byte_mv, cursor, _dtype, itemsize, value) -> int Load an object into the numpy byte view at index cursor, return the index of the end of the object :param byte_mv: numpy byte view :param cursor: index of where the object start :param _dtype: data type of the object :param itemsize: size of the data type :param value: value to write :returns: end of object index .. py:function:: mv_write_summary_header(byte_mv, cursor, event_id, summary_id, exposure_value) -> int Wrapper for cached write a summary header to the numpy byte view at index cursor, return the index of the end of the object :param byte_mv: numpy byte view :param cursor: index of where the object start :param event_id: event id :param summary_id: summary id :param exposure_value: exposure value :returns: end of object index .. py:function:: mv_write_summary_header_cached(byte_mv, cursor, event_id, summary_id, exposure_value, event_id_type, event_id_size, summary_id_dtype, summary_id_size, exposure_value_dtype, exposure_value_size) -> int Cached write a summary header to the numpy byte view at index cursor, return the index of the end of the object :param byte_mv: numpy byte view :param cursor: index of where the object start :param event_id: event id :param summary_id: summary id :param exposure_value: exposure value :param event_id_type: type info for event id :param event_id_size: size of event id in bytes :param summary_id_dtype: type info for summary id :param summary_id_size: size of summary id :param exposure_value_dtype: type info for exposure value :param exposure_value_size: size of exposure value :returns: end of object index .. py:function:: mv_write_item_header(byte_mv, cursor, event_id, item_id) -> int Wrapper function for cached mv_write_item_header. writes an item header to the numpy byte view, return index of the end of the object :param byte_mv: numpy byte view :param cursor: index of where the object start :param event_id: event id :param item_id: item id :returns: end of object index .. py:function:: mv_write_sidx_loss(byte_mv, cursor, sidx, loss) -> int Write sidx and loss to the numpy byte view at index cursor, return the index of the end of the object :param byte_mv: numpy byte view :param cursor: index of where the object start :param sidx: sample id :param loss: loss :returns: end of object index .. py:function:: mv_write_sidx_loss_cached(byte_mv, cursor, sidx, loss, sidx_type, loss_type, sidx_size, loss_size) -> int Cached write sidx and loss to the numpy byte view at index cursor, return the index of the end of the object :param byte_mv: numpy byte view :param cursor: index of where the object start :param sidx: sample id :param loss: loss :param sidx_type: sample id type info :param loss_type: loss type info :param sidx_size: sample id size in bytes :param loss_size: loss size in bytes :returns: end of object index .. py:class:: EventReader Abstract class to read event stream This class provides a generic interface to read multiple event streams using: - **selector**: handle back pressure — the program is paused and doesn't use resources if nothing is in the stream buffer - **memoryview**: read a chunk (PIPE_CAPACITY) of data at a time then work on it using a numpy byte view of this buffer These methods need to be implemented: - ``__init__(self, ...)``: the constructor with all data structures needed to read and store the event stream - ``read_buffer(self, byte_mv, cursor, valid_buff, event_id, item_id)``: point to a local numba.jit function named read_buffer (a template is provided below); it should implement the specific logic of where and how to store the event information These methods may be overwritten: - ``item_exit(self)``: specific logic to do when an item is finished (only executed once the stream is finished but no 0,0 closure was present) - ``event_read_log(self)``: what kpi to log when a full event is read Usage snippet:: with ExitStack() as stack: streams_in, (stream_type, stream_agg_type, len_sample) = init_streams_in(files_in, stack) reader = CustomReader() for event_id in reader.read_streams(streams_in): .. py:method:: register_streams_in(selector_class, streams_in) :staticmethod: Data from input process is generally sent by event block, meaning once a stream receive data, the complete event is going to be sent in a short amount of time. Therefore, we can focus on each stream one by one using their specific selector 'stream_selector'. .. py:method:: read_streams(streams_in) Read multiple stream input, yield each event id and load relevant value according to the specific read_buffer implemented in subclass :param streams_in: streams to read :Yields: *int* -- each event id read from the streams .. py:method:: read_event(stream_in, main_selector, stream_selector, mv, byte_mv, cursor, valid_buff, file_idx) Read one event from stream_in close and remove the stream from main_selector when all is read :param stream_in: stream to read :param main_selector: selector that contain all the streams :param stream_selector: this stream selector :param mv: buffer memoryview :param byte_mv: numpy byte view of the buffer :param cursor: current cursor of the memory view :param valid_buff: valid data in memory view :param file_idx: file index :returns: event_id, cursor, valid_buff .. py:method:: read_buffer(byte_mv, cursor, valid_buff, event_id, item_id, **kwargs) :abstractmethod: .. py:method:: item_exit() .. py:method:: finalize_event() Hook called by read_streams just before yielding each event_id. Subclasses use this to perform any per-event finishing work that read_buffer cannot do on its own — e.g. sorting accumulated state. Default is a no-op. .. py:method:: event_read_log(event_id) .. py:function:: write_mv_to_stream(stream, byte_mv, cursor) Write numpy byte array view to stream - use select to handle forward pressure - use a while loop in case the stream is non-blocking (meaning the ammount of byte written is not guarantied to be cursor len) :param stream: stream to write to :param byte_mv: numpy byte view of the buffer to write :param cursor: ammount of byte to write