Pickles `message` instances, stores it in the provided `shm` chunk at given offset and returns `ShmMessageDesc` instance describing the placement of the `message`. Returned instance can be put into ShmQueue.
(worker_id, shm_chunk: BufShmChunk, message, offset, resize=True)
| 297 | |
| 298 | |
| 299 | def write_shm_message(worker_id, shm_chunk: BufShmChunk, message, offset, resize=True): |
| 300 | """ |
| 301 | Pickles `message` instances, stores it in the provided `shm` chunk at given offset and returns |
| 302 | `ShmMessageDesc` instance describing the placement of the `message`. |
| 303 | Returned instance can be put into ShmQueue. |
| 304 | """ |
| 305 | serialized_message = pickle.dumps(message) |
| 306 | num_bytes = len(serialized_message) |
| 307 | if num_bytes > shm_chunk.capacity - offset: |
| 308 | if resize: |
| 309 | resize_shm_chunk(shm_chunk, offset + num_bytes) |
| 310 | else: |
| 311 | # This should not happen, resize is False only when writing task description into memory |
| 312 | # in the main process, and the description (ScheduledTask and its members) boils down |
| 313 | # to bounded number of integers. |
| 314 | raise RuntimeError( |
| 315 | "Could not put message into shared memory region," |
| 316 | " not enough space in the buffer." |
| 317 | ) |
| 318 | buffer = shm_chunk.buf[offset : offset + num_bytes] |
| 319 | buffer[:] = serialized_message |
| 320 | return ShmMessageDesc(worker_id, shm_chunk.shm_chunk_id, shm_chunk.capacity, offset, num_bytes) |
no test coverage detected