MCPcopy Create free account
hub / github.com/NVIDIA/DALI / ShmQueue

Class ShmQueue

dali/python/nvidia/dali/_multiproc/shared_queue.py:27–182  ·  view source on GitHub ↗

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.

Source from the content-addressed store, hash-verified

25
26
27class 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

Callers 4

create_worker_contextsFunction · 0.90
from_contextsMethod · 0.90
setup_queue_and_workerFunction · 0.90

Calls

no outgoing calls

Tested by 2

setup_queue_and_workerFunction · 0.72