Simple fixed capacity shared memory queue of fixed size messages. Writing to a full queue fails, attempt to get from an empty queue blocks until data is available or the queue is closed.
| 25 | |
| 26 | |
| 27 | class ShmQueue: |
| 28 | """ |
| 29 | Simple fixed capacity shared memory queue of fixed size messages. |
| 30 | Writing to a full queue fails, attempt to get from an empty queue blocks until data is |
| 31 | available or the queue is closed. |
| 32 | """ |
| 33 | |
| 34 | MSG_CLASS = ShmMessageDesc |
| 35 | ALIGN_UP_MSG = 4 |
| 36 | ALIGN_UP_BUFFER = 4096 |
| 37 | |
| 38 | def __init__(self, mp, capacity): |
| 39 | self.lock = mp.Lock() |
| 40 | self.cv_not_empty = mp.Condition(self.lock) |
| 41 | self.capacity = capacity |
| 42 | self.meta = QueueMeta(capacity, 0, 0, 0) |
| 43 | self.meta_size = align_up(self.meta.get_size(), self.ALIGN_UP_MSG) |
| 44 | dummy_msg = self.MSG_CLASS() |
| 45 | self.msg_size = align_up(dummy_msg.get_size(), self.ALIGN_UP_MSG) |
| 46 | self.shm_capacity = align_up( |
| 47 | self.meta_size + capacity * self.msg_size, self.ALIGN_UP_BUFFER |
| 48 | ) |
| 49 | self.shm = shared_mem.SharedMem.allocate(self.shm_capacity) |
| 50 | self.is_closed = False |
| 51 | self._init_offsets() |
| 52 | self._write_meta() |
| 53 | |
| 54 | def __getstate__(self): |
| 55 | state = self.__dict__.copy() |
| 56 | state["msgs_offsets"] = None |
| 57 | state["shm"] = None |
| 58 | return state |
| 59 | |
| 60 | def __setstate__(self, state): |
| 61 | self.__dict__.update(state) |
| 62 | self._init_offsets() |
| 63 | |
| 64 | def _init_offsets(self): |
| 65 | self.msgs_offsets = [i * self.msg_size + self.meta_size for i in range(self.capacity)] |
| 66 | |
| 67 | def _read_meta(self): |
| 68 | self.meta.unpack_from(self.shm.buf, 0) |
| 69 | |
| 70 | def _write_meta(self): |
| 71 | self.meta.pack_into(self.shm.buf, 0) |
| 72 | |
| 73 | def _read_msg(self, i): |
| 74 | offset = self.msgs_offsets[i] |
| 75 | msg = self.MSG_CLASS() |
| 76 | msg.unpack_from(self.shm.buf, offset) |
| 77 | return msg |
| 78 | |
| 79 | def _write_msg(self, i, msg): |
| 80 | offset = self.msgs_offsets[i] |
| 81 | msg.pack_into(self.shm.buf, offset) |
| 82 | |
| 83 | def _recv_samples(self, num_samples): |
| 84 | num_take = self.meta.size |
no outgoing calls