ipcq 0.12.0
Loading...
Searching...
No Matches
shm.hpp
Go to the documentation of this file.
1/**
2 * @file
3 * @ingroup ipcq_detail
4 * @brief Contains core shm implementation details.
5 * @copyright
6 * SPDX-FileCopyrightText: 2018-2024 European Southern Observatory (ESO)
7 *
8 * SPDX-License-Identifier: LGPL-3.0-only
9 */
10#ifndef IPCQ_DETAIL_SHM_HPP_
11#define IPCQ_DETAIL_SHM_HPP_
12#include <atomic>
13#include <array>
14#include <chrono>
15#include <functional>
16#include <signal.h> // kill
17#include <system_error>
18#include <thread>
19#include <type_traits>
20
21#include <boost/version.hpp>
22// Headers for BoostIpcCondition
23#include <boost/date_time/microsec_time_clock.hpp>
24#include <boost/date_time/posix_time/posix_time.hpp>
25#include <boost/interprocess/mapped_region.hpp>
26#include <boost/interprocess/shared_memory_object.hpp>
27#include <boost/interprocess/sync/interprocess_condition.hpp>
28#include <boost/interprocess/sync/interprocess_mutex.hpp>
29#include <boost/interprocess/sync/scoped_lock.hpp>
30
31#include <numapp/mempolicy.hpp>
32
33#include "../error.hpp"
34#include "iterator.hpp"
35#include "traits.hpp"
36
37#if __GNUC__
38#define restrict __restrict__
39#else
40#define restrict
41#endif
42
43#ifndef NDEBUG
44#define IPCQ_ENABLE_ASSERT
45#define IPCQ_ENABLE_LOG
46#endif
47
48#ifdef IPCQ_ENABLE_ASSERT
49#define IPCQ_ASSERT(x) \
50 do { \
51 assert(x); \
52 } while (0)
53#else
54#define IPCQ_ASSERT(x)
55#endif
56
57#ifdef IPCQ_ENABLE_LOG
58#define IPCQ_LOG(x) \
59 do { \
60 std::cout << x << std::endl; \
61 } while (0)
62#else
63#define IPCQ_LOG(x)
64#endif
65
66/** Version number */
67#define IPCQ_FILE_VERSION 1
68
69namespace ipcq::detail {
70namespace config {
71
72/**
73 * File magic number that identifies ipcq shared objects.
74 */
75constexpr std::array<uint8_t, 7> MAGIC_NUMBER = {0x69, 0x70, 0x63, 0x71, 0xF0, 0x6A, 0x9C};
76
77/**
78 * Prefix used for shared object names.
79 */
80constexpr std::string_view SHM_OBJECT_NAME_PREFIX = "ipcq-";
81
82/**
83 * Number of bits making up first, last and size each for a total of 63 bits
84 *
85 * @ingroup ipcq_detail
86 */
87constexpr auto BIT_COUNT = 21u;
88/**
89 * Bitmask for bit count
90 *
91 * @ingroup ipcq_detail
92 */
93constexpr auto BIT_MASK = 0x1f'ffffu;
94
95/**
96 * Maximum number of elements in queue.
97 * `(2 ^ BIT_COUNT) - 1`
98 * @ingroup ipcq_detail
99 */
100constexpr auto MAX_VALUE = 2097151u;
101
102/**
103 * Bit indicating if queue is closed (highest bit).
104 *
105 * @ingroup ipcq_detail
106 */
107constexpr auto BIT_IS_CLOSED = 0x8000'0000'0000'0000u;
108
109/**
110 * Length of typename character sequences.
111 *
112 * @ingroup ipcq_detail
113 */
114constexpr size_t TYPENAME_SEQ_LEN = 128;
115
116} // namespace config
117
118namespace ipc = boost::interprocess;
119
120/**
121 * Static parameters for a given shm queue.
122 */
123struct QueueInfo {
124 static constexpr QueueInfo Make(size_t shm_size,
125 size_t capacity,
127 size_t counter_offset,
128 size_t element_offset,
129 size_t element_size) noexcept {
130 return QueueInfo{
132 }
133
134 /**
135 * Number of bytes.
136 */
137 size_t shm_size;
138 /**
139 * Queue capacity in number of elements.
140 */
141 size_t capacity;
143 /**
144 * Offset to std::atomic<CounterType> from &Element<T>.
145 */
147 /**
148 * Offset to std::atomic<CounterType> from &Element<T>::element.
149 */
151 /**
152 * Size of Element<T>.
153 */
155};
156
157static_assert(std::is_trivially_copyable_v<QueueInfo>, "must be trivial since it resides in SHM");
158
159/**
160 * Prefault memory by reading a byte in every page starting at @a start up to @a lenght bytes.
161 *
162 * To avoid optimizing this no-op the pointer is declared volatile.
163 */
164inline void PrefaultPages(void const* begin, void const* end) noexcept {
165 volatile const uint8_t* ptr = static_cast<uint8_t const*>(begin);
166 size_t pz = getpagesize();
167 for (; ptr < end; ptr += pz) {
168 (void)*ptr;
169 }
170}
171
172/**
173 * Control block residing in shared memory shared by readers and writers.
174 * The control block only uses offsets since base address for each user is different.
175 *
176 * @ingroup ipcq_detail
177 */
179 /** Magic number that identifies ipcq files*/
180 std::array<uint8_t, 7> magic;
181 /** IPCQ protocol version number (not IPCQ version number)*/
182 uint8_t version;
183
184 // They are frequently updated together (but not atomically).
185 /** Range contains the compressed information about begin, end of queue. */
186 std::atomic<unsigned long long> range;
187 size_t element_type_hash; ///< std::hash of element type T
188 size_t condition_type_hash; ///< std::hash of ConditionPolicy
190 /**
191 * Non-zero process id of owner if owner is available or
192 * 0 if owner gracefully destructed itself but could
193 * not free shared memory because it is still allocated by readers.
194 */
195 std::atomic<pid_t> owner_pid;
196 /** @name Parameters used by ipcq-spy */
197 // @{
198 // Char sequences may or may not be terminated with \0.
201 // @}
202};
203
204static_assert(std::is_trivially_copyable<ShmControlBlock>::value,
205 "ShmControlBlock must be a POD type since it located in shared memory");
206
207static_assert(sizeof(ShmControlBlock) % 8 == 0);
208
209/**
210 * Helper that adds to posix_time::ptime a std:chrono duration.
211 */
212template <class Rep, class Period>
213boost::posix_time::ptime AddDuration(boost::posix_time::ptime const& point,
214 std::chrono::duration<Rep, Period> duration) noexcept {
215 using namespace std::chrono;
216 using usec = std::chrono::duration<int64_t, std::micro>; // because std only specifies 55 bits
217 return point + boost::posix_time::microseconds(duration_cast<usec>(duration).count());
218}
219
220/**
221 * Creates the shared memory object name from topic_name.
222 *
223 * @param topic_name The ipcq shared memory topic name.
224 */
225inline std::string MakeShmObjectName(char const* topic_name) {
226 std::string base(config::SHM_OBJECT_NAME_PREFIX);
227 base += topic_name;
228 return base;
229}
230
231/**
232 * Small helper that removes shared memory, to make ShmControl exception-safe with RAII.
233 *
234 * @ingroup ipcq_detail
235 */
236template <class ShmTraits>
238public:
239 using SharedMemoryObject = typename ShmTraits::SharedMemoryObject;
240 using MappedRegion = typename ShmTraits::MappedRegion;
241
242 enum class Policy { None = 0, Destructor = 1, ConstructorAndDestructor = 2 };
243 ShmRemover(std::string name, Policy policy) : m_shm_name(std::move(name)), m_policy(policy) {
244 if (m_policy == Policy::ConstructorAndDestructor) {
245 SharedMemoryObject::remove(m_shm_name.c_str());
246 }
247 }
248 ShmRemover(ShmRemover&& rhs) noexcept
249 : m_shm_name(std::move(rhs.m_shm_name)), m_policy(rhs.m_policy) {
250 // Moved from object should do nothing.
251 rhs.m_shm_name = "";
252 rhs.m_policy = Policy::None;
253 }
254 ShmRemover& operator=(ShmRemover&& rhs) noexcept {
255 m_shm_name = std::move(rhs.m_shm_name);
256 m_policy = rhs.m_policy;
257 // Moved from object should do nothing.
258 rhs.m_shm_name = "";
259 rhs.m_policy = Policy::None;
260 return *this;
261 }
262 ~ShmRemover() noexcept {
263 if (m_policy == Policy::ConstructorAndDestructor || m_policy == Policy::Destructor) {
264 (void)SharedMemoryObject::remove(m_shm_name.c_str());
265 }
266 }
267
268private:
269 std::string m_shm_name;
270 Policy m_policy;
271};
272
273/**
274 * Simple spinning condition
275 *
276 * @ingroup ipcq_detail
277 */
279public:
280 using condition_t = unsigned int;
281
282 void Signal() noexcept {
283 m_condition.fetch_add(1, std::memory_order_release);
284 }
285
286 /**
287 * @note: It is possible that between detecting that a wait is necessary and actually creating
288 * the initial state to await a change on, the writer calls Signal();
289 * In this scenario the reader will not be able to see the signal as a true signal.
290 *
291 * @todo: Mitigate this data race by reading range inside wait?
292 */
293 template <class Rep, class Period, class Predicate>
294 [[nodiscard]] bool Wait(std::chrono::duration<Rep, Period> timeout, Predicate pred) noexcept {
295 auto start = std::chrono::high_resolution_clock::now();
296 while (!pred()) {
297 auto initial = m_condition.load(std::memory_order_acquire);
298 // Spin until condition variable changes, indicating a notication or that we time out.
299 while (initial == m_condition.load(std::memory_order_acquire)) {
300 if ((std::chrono::high_resolution_clock::now() - start) > timeout) {
301 return false;
302 }
303 std::this_thread::yield();
304 }
305 }
306 return true;
307 }
308
309private:
310 std::atomic<unsigned int> m_condition;
311};
312
313/**
314 * Condition policy based on boost::interprocess::interprocess_condition.
315 *
316 * @ingroup ipcq_detail
317 */
319public:
320 // Needs a default constructor
322 /**
323 * Unlike std::condition_variable in C++11, it is NOT safe to invoke the destructor if all
324 * threads have been only notified. It is required that they have exited their respective wait
325 * functions.
326 */
328 m_cond.notify_all();
329 // Sleep for a while to let any reader wake up?
330 }
331
332 void Signal() noexcept {
333 {
334 // Prevent signal while readers hold lock so that the window between running predicate
335 // and calling wait on m_cond is removed (which would make it possible to miss
336 // a notification).
337 ipc::scoped_lock<ipc::interprocess_mutex> lock(m_mtx);
338 }
339 m_cond.notify_all();
340 }
341
342 template <class Rep, class Period, class Predicate>
343 [[nodiscard]] bool Wait(std::chrono::duration<Rep, Period> timeout, Predicate pred) noexcept {
344 ipc::scoped_lock<ipc::interprocess_mutex> lock(m_mtx);
345 return m_cond.timed_wait(
346 lock,
347 AddDuration(boost::posix_time::microsec_clock::universal_time(), timeout),
348 std::forward<Predicate>(pred));
349 }
350
351private:
352 ipc::interprocess_mutex m_mtx;
353 ipc::interprocess_condition m_cond;
354};
355
356// @implements ShmTrait used by ShmControl
358 using SharedMemoryObject = boost::interprocess::shared_memory_object;
359 using MappedRegion = boost::interprocess::mapped_region;
360};
361
362/**
363 * Non-virtual base-class for ShmControl that does not depend on any static type.
364 *
365 * This is mainly a requirement to support ipcq-spy which is not templated on element type or
366 * any policies.
367 *
368 * @tparam ShmTraits Trait type used to facilitate testing. Does not depend on user provided types.
369 */
370template <class ShmTraits = BoostInterprocessTraits>
372public:
373 using CounterType = size_t;
374 using SharedMemoryObject = typename ShmTraits::SharedMemoryObject;
375 using MappedRegion = typename ShmTraits::MappedRegion;
376
377 /**
378 * Owner constructor.
379 *
380 * @param size Data buffer in number of bytes.
381 * @param data_offset Offset to first byte of the contigous queue buffer, relative m_queue
382 * region.
383 */
384 explicit ShmControlBase(char const* name,
385 QueueInfo queue_info,
386 std::optional<numapp::MemPolicy> mem_policy = std::nullopt,
387 int mmap_flags = 0)
388 : m_topic_name(name)
389 , m_queue_info(queue_info)
390 , m_shm_obj{ipc::create_only, MakeShmObjectName(name).c_str(), ipc::read_write}
391 , m_shmem_remover{m_shm_obj.get_name(), ShmRemover<ShmTraits>::Policy::Destructor}
392 , m_ctl_block{nullptr}
393 , m_is_owner(true) {
394 if (m_queue_info.capacity > config::MAX_VALUE) {
395 throw std::system_error(make_error_code(Error::Overflow),
396 "Requested capacity overflows storage capacity");
397 }
398 // Normally mbind(2) could be used to specify the memory policy when using mmap(2). But as
399 // it says in mbind(2):
400 //
401 // The specified policy will be ignored for any MAP_SHARED mappings in
402 // the specified memory range. Rather the pages will be allocated
403 // according to the memory policy of the thread that caused the page to
404 // be allocated. Again, this may not be the thread that called mbind().
405 //
406 // instead we would have to:
407 // 1. Allocate (and not touch memory!)
408 // 2. Set default mem policy (not with `mbind(2)`)
409 // 3. Set mlock to avoid swapping out getting different policy next time.
410 // 4. Fault all memory.
411 //
412 // And that is what is essentially done here by setting default memory policy for this
413 // thread and then making sure that all memory is paged in (with MAP_POPULATE) but
414 // with two sys-calls (`set_mempolicy(2)` and then `mmap(2)`). Three if we count the second
415 // `set_mempolicy` to restore the previous policy after shared memory allocation.
416 //
417 std::optional<numapp::ScopedMemPolicy> scoped_policy;
418 if (mem_policy) {
419 scoped_policy.emplace(*mem_policy);
420 }
421 // Sets the size of the shared memory object to shm_size
422#if BOOST_VERSION >= 108500
423 m_shm_obj.truncate(m_queue_info.shm_size);
424#else
425 // Workaround for missing error code from `posix_fallocate()` fixed in boost 1.85
426 // https://github.com/boostorg/interprocess/issues/166
427 try {
428 errno = 0;
429 m_shm_obj.truncate(m_queue_info.shm_size);
430 } catch (boost::interprocess::interprocess_exception const& ex) {
431 if (ex.get_native_error() == 0) {
432 throw boost::interprocess::interprocess_exception(
433 boost::interprocess::error_info(ENOSPC));
434 } else {
435 throw;
436 }
437 }
438#endif
439
440 m_control_region = std::move(MappedRegion(
441 m_shm_obj,
442 ipc::read_write,
443 0, // control begins at 0
444 sizeof(ShmControlBlock), // size of region
445 nullptr,
446 (::boost::interprocess::default_map_options | MAP_LOCKED | MAP_POPULATE | mmap_flags)));
448 std::move(MappedRegion(m_shm_obj,
449 ipc::read_write,
450 sizeof(ShmControlBlock), // begins where m_control_region ends
451 0, // size of region is the rest (0)
452 nullptr,
453 ipc::default_map_options | MAP_LOCKED | MAP_POPULATE));
454
455 // Initialize control block in shared memory.
456 m_ctl_block = static_cast<ShmControlBlock*>(m_control_region.get_address());
459 std::memcpy(&m_ctl_block->queue_spec,
460 &queue_info,
461 sizeof(queue_info));
462 std::memcpy(&m_queue_info,
463 &queue_info,
464 sizeof(m_queue_info));
465 }
466
467 /**
468 * Non-owner constructor.
469 */
470 explicit ShmControlBase(char const* name)
471 : m_topic_name(name)
472 , m_shm_obj{ipc::open_only,
473 (std::string(config::SHM_OBJECT_NAME_PREFIX) + name).c_str(),
474 ipc::read_write}
475 , m_shmem_remover{m_shm_obj.get_name(), ShmRemover<ShmTraits>::Policy::None}
476 , m_is_owner{false} {
477 // We don't want readers to destruct the shared memory in constructor
479 std::move(MappedRegion(m_shm_obj,
480 ipc::read_only,
481 0, // control begins at 0
482 sizeof(ShmControlBlock), // size
483 nullptr, // addr
484 ipc::default_map_options | MAP_LOCKED | MAP_POPULATE));
485 m_queue_region = std::move(
487 ipc::read_write, // rw needed for condition policy
488 sizeof(ShmControlBlock), // start at offset by size of control block
489 0,
490 nullptr, // addr
491 ipc::default_map_options | MAP_LOCKED | MAP_POPULATE));
492
493 m_ctl_block = static_cast<ShmControlBlock*>(m_control_region.get_address());
494 if (this->m_ctl_block->magic != config::MAGIC_NUMBER) {
495 throw std::system_error(make_error_code(Error::BadFile),
496 m_shm_obj.get_name());
497 }
498 std::memcpy(&m_queue_info,
499 &m_ctl_block->queue_spec,
500 sizeof(m_queue_info));
501 }
502
504 : m_topic_name(std::move(rhs.m_topic_name))
505 , m_queue_info(rhs.m_queue_info)
506 , m_shm_obj(std::move(rhs.m_shm_obj))
507 , m_shmem_remover(std::move(rhs.m_shmem_remover))
508 , m_control_region(std::move(rhs.m_control_region))
509 , m_queue_region(std::move(rhs.m_queue_region))
510 , m_ctl_block(rhs.m_ctl_block)
511 , m_is_owner(rhs.m_is_owner) {
512 // rhs no longer own the object.
513 rhs.m_is_owner = false;
514 rhs.m_queue_info.shm_size = 0u;
515 }
516
518 m_topic_name = std::move(rhs.m_topic_name);
519 m_queue_info = rhs.m_queue_info;
520 m_shm_obj = std::move(rhs.m_shm_obj);
521 m_shmem_remover = std::move(rhs.m_shmem_remover);
522 m_control_region = std::move(rhs.m_control_region);
523 m_queue_region = std::move(rhs.m_queue_region);
524 m_ctl_block = rhs.m_ctl_block;
525 m_is_owner = rhs.m_is_owner;
526 // rhs no longer own the object.
527 rhs.m_is_owner = false;
528 rhs.m_queue_info.shm_size = 0u;
529 return *this;
530 }
531
532 [[nodiscard]] char const* TopicName() const noexcept {
533 return m_topic_name.c_str();
534 }
535
536 [[nodiscard]] char const* ShmObjectName() const {
537 return m_shm_obj.get_name();
538 }
539
540 [[nodiscard]] constexpr size_t Capacity() const noexcept {
541 return m_queue_info.capacity;
542 }
543
544 /**
545 * @returns ipcq shared object file/protocol version
546 */
547 [[nodiscard]] constexpr uint8_t FileVersion() const noexcept {
548 return this->m_ctl_block->version;
549 }
550
551 /**
552 * @returns PID of registered owner (writer) process.
553 */
554 [[nodiscard]] constexpr pid_t OwnerPid() const noexcept {
555 return this->m_ctl_block->owner_pid;
556 }
557
558protected:
559 /// Topic name, as provided by user.
560 std::string m_topic_name;
562
563 /// Manages the shared memory
565 // Helper object that deletes the allocated shared memory
567 /// Mapped segment to the control (mapped as read only for Reader).
569 /// Mapped region to the queue (mapped as read only for Reader).
571 /// @name Pointers to mapped shared memory regions
572 // @{
574 // @}
575
576 bool m_is_owner; ///< Indicates if this object owns the shared memory or not
577};
578
579/**
580 * ShmControl provides the basic data model for the queue and low level operations
581 * on the shared memory queue, some of which are atomic.
582 *
583 * The shared data model for the queue includes:
584 *
585 * - a Range describing the valid range of the data and
586 * - the capacity of the circular buffer.
587 *
588 * This information can be compressed into a single unsigned long long value
589 * which can be atomically read or writen.
590 *
591 * The atomic operations are:
592 *
593 * - `AtomicLoadState()'
594 * - `AtomicStoreState()`
595 *
596 *
597 * ShmControl does not encode knowledge of the buffer semantics. It merely provides the shared
598 * access patterns to it.
599 *
600 * @todo: Optimize atomic stores. Since x86 (and AMD64) do not reorder stores with other stores and
601 * loads with other loads a single operation could be used to load both range and generation as well
602 * as store range and generation.
603 *
604 * Two mappings are created to the shm segment, one for the control block and one for the queue.
605 * It's not really important, but if each region is mapped individually, there's less chance that an
606 * accidental overwrite of the queue region affects the control region. It also make offset
607 * calculations simpler since the control block is not accounted for.
608 *
609 * Writer Protocol
610 * ---------------
611 *
612 * Writers to the queue should adhere to the following protocol to ensure data consistency or detect
613 * inconsistent states:
614 *
615 * To write one or more elements:
616 *
617 * 1. Reserve enough room for elements to be written by calling `Pop()` once for each element to
618 * write.
619 * 2. "Publish" change of queue state by calling `AtomicStoreState()`.
620 * 3. Add elements by calling `Push()`.
621 * 4. "Publish" change of queue again by calling `AtomicStoreState()`.
622 *
623 * Reader Protocol
624 * ---------------
625 *
626 * 1. Read current queue state with `AtomicLoadState()`
627 * 2. Verify that iterator is valid.
628 * 3. Copy.
629 * 4. Read current queue state with `AtomicLoadState()`
630 * 5. Verify that iterator is still valid.
631 * 6. Verify that copied element has the expected counter value.
632 *
633 * Important
634 * ---------
635 *
636 * boost interprocess uses mmap on linux, and it is possible that write combine is used
637 * under unknown circumstances. If write combining is enabled then the strong
638 * memory ordering is no longer true.
639 *
640 * Implementation Notes
641 * --------------------
642 *
643 * boost::shared_memory_object corresponds to a file handle, opened with `open(2)`.
644 * boost::shared_memory_object::truncate corresponds to `ftruncate(2)`.
645 *
646 * boost::mapped_region construction corresponds to a call to `mmap(2)` which in
647 * the case of ShmControl also allocate, lock and populate the memory using
648 * MAP_LOCKED and MAP_POPULATE.
649 *
650 * User control over memory policy for the allocation is made possible because
651 * ShmControl in the "writer" is always the one that creates the boost::shared_memory_object
652 * and can then make sure that the chosen memory policy is applied when the allocation
653 * is performed in `mmap(2)`.
654 *
655 * @ingroup ipcq_detail
656 */
657template <class T, class ConditionPolicy, class ShmTraits = BoostInterprocessTraits>
658class ShmControl final : public ShmControlBase<ShmTraits> {
659public:
661 using SharedMemoryObject = typename ShmTraits::SharedMemoryObject;
662 struct Element {
663 std::atomic<CounterType> counter;
665 };
666
671 using difference_type = ptrdiff_t;
674
675 friend Iterator;
677
678 // Object in shared memory
679 struct Queue {
680 ConditionPolicy condition;
681 value_type buffer[1]; // beginning of buffer
682 };
683
684 /**
685 * Attempts to remove shared memory object.
686 *
687 * @note @a shm_object_name is not the topic name, but the raw object name.
688 *
689 * @param shm_object_name The raw object name to remove.
690 * @returns false on error.
691 */
692 static std::error_code RemoveShmObject(char const* shm_object_name) noexcept {
693 errno = 0;
694 // Returns false on error
695 if (!SharedMemoryObject::remove(shm_object_name)) {
696 return std::make_error_code(std::errc(errno));
697 }
698 return {};
699 }
700
701 /**
702 * Create a new (owned) named shared memory segment.
703 *
704 * @warning Will (try to) destroy shared memory with same name if it exists.
705 *
706 * @param name Full name of shared memory segment.
707 * @param capacity Capacity of queue. Maximum capacity is config::MAX_CAPACITY.
708 * @param mem_policy Optional memory policy that is applied when allocating the shared memory.
709 * @param mmap_flags Additional flags that will be or-ed with the default flags
710 * MAP_LOCKED and MAP_POPULATE.
711 *
712 * @throws boost::interprocess::interprocess_exception if creating shared memory fails.
713 * @throws std::system_error with
714 * - Error::Overflow if requested capacity is greater than what can be represented by the
715 * queue.
716 */
717 explicit ShmControl(char const* name,
718 size_t capacity,
719 std::optional<numapp::MemPolicy> mem_policy = std::nullopt,
720 int mmap_flags = 0)
721 : ShmControlBase<ShmTraits>(
722 name,
723 QueueInfo::Make(
724 sizeof(ShmControlBlock) + offsetof(Queue, buffer) + sizeof(value_type) * capacity,
725 capacity,
726 offsetof(Queue, buffer),
727 offsetof(Element, counter),
728 offsetof(Element, element),
729 sizeof(Element)
730 ),
731 mem_policy,
732 mmap_flags)
733 , m_buffer{nullptr}
734 , m_capacity{capacity}
735 , m_first{nullptr}
736 , m_last{nullptr}
737 , m_size{0u}
738 , m_is_closed{false}
739 , m_counter{0} {
740 assert(this->m_ctl_block);
741
742 m_queue = new (this->m_queue_region.get_address()) Queue();
743
744 m_buffer = m_queue->buffer;
745 m_first = m_buffer;
746 m_last = m_buffer;
747
748 {
749 // Possibly truncated string is expected and allowed.
750 #pragma GCC diagnostic push
751 #pragma GCC diagnostic ignored "-Wpragmas"
752 #pragma GCC diagnostic ignored "-Wstringop-truncation"
753 this->m_ctl_block->element_type_hash = std::hash<std::string>{}(typeid(T).name());
754 strncpy(this->m_ctl_block->element_type_name, typeid(T).name(),
756 this->m_ctl_block->condition_type_hash =
757 std::hash<std::string>{}(typeid(ConditionPolicy).name());
758 strncpy(this->m_ctl_block->condition_type_name, typeid(ConditionPolicy).name(),
760 #pragma GCC diagnostic pop
762 }
763 // Finally mark the queue available by setting the pid
764 this->m_ctl_block->owner_pid.store(getpid(), std::memory_order_release);
765 }
766
767 /**
768 * Open an existing (not owned) shared memory segment in read only mode.
769 *
770 * @param name Full name of shared memory segment.
771 *
772 * @throws std::system_error If
773 * - attaching to shared memory fail
774 * - ipcq::Error::TypeMismatch if element or condition policy types don't match
775 * - ipcq::Error::Closed if queue is closed or owner does not exist.
776 */
777 explicit ShmControl(char const* name)
778 : ShmControlBase<ShmTraits>(name)
779 , m_first{nullptr}
780 , m_last{nullptr}
781 , m_size{0u}
782 , m_is_closed{false}
783 , m_counter{0} {
784 assert(this->m_ctl_block);
785
786 // Check if writer exists by checking the pid, which is the last synchronization being done
787 // by writer before queue is usable. Therefore this is checked first.
788 //
789 // Writer will set pid to zero if it gracefully exits, on crashes it may still be set
790 // but process is gone, so we check that the process exist with `kill()`.
791 auto owner_pid = this->m_ctl_block->owner_pid.load(std::memory_order_acquire);
792 auto owner_exists = owner_pid != 0 && kill(owner_pid, 0) == 0;
793
794 if (this->m_ctl_block->version != IPCQ_FILE_VERSION) {
795 throw std::system_error(make_error_code(Error::BadVersion),
796 this->ShmObjectName());
797 }
798
799 if (!owner_exists) {
800 throw std::system_error(make_error_code(Error::Closed), "Writer does not exist");
801 }
802
803 if (this->m_ctl_block->element_type_hash != std::hash<std::string>{}(typeid(T).name())) {
804 throw std::system_error(make_error_code(Error::TypeMismatch), "Element type mismatch");
805 }
806
807 if (this->m_ctl_block->condition_type_hash !=
808 std::hash<std::string>{}(typeid(ConditionPolicy).name())) {
809 throw std::system_error(make_error_code(Error::TypeMismatch),
810 "ConditionPolicy type mismatch");
811 }
812
813 m_queue = static_cast<Queue*>(this->m_queue_region.get_address());
814 m_buffer = m_queue->buffer;
815 m_capacity = this->Capacity();
816
817 // Check if writer exists by checking the pid.
818 // writer will set it to zero if it gracefully exits, on crashes it may still be set
819 // but process is gone, so we check that as well.
821 if (m_is_closed) {
822 throw std::system_error(make_error_code(Error::Closed),
823 "Writer exists but queue is closed");
824 }
825 }
826
827 ~ShmControl() noexcept {
828 if (this->m_is_owner) {
829 Close();
831 this->m_ctl_block->owner_pid = 0;
832 m_queue->~Queue();
833 }
834 }
835
836 /**
837 * Moves from rhs into this.
838 *
839 * @pre rhs is valid.
840 * @post rhs is invalid.
841 */
842 ShmControl(ShmControl&& rhs) noexcept
843 : ShmControlBase<ShmTraits>(std::move(rhs))
844 , m_queue(rhs.m_queue)
845 , m_buffer(rhs.m_buffer)
846 , m_capacity(rhs.m_capacity)
847 , m_first(rhs.m_first)
848 , m_last(rhs.m_last)
849 , m_size(rhs.m_size)
850 , m_is_closed(rhs.m_is_closed)
851 , m_counter(rhs.m_counter) {
852 // rhs no longer own the object.
853 rhs.m_is_owner = false;
854 rhs.m_queue = nullptr;
855 rhs.m_buffer = nullptr;
856 rhs.m_first = nullptr;
857 rhs.m_last = nullptr;
858 rhs.m_size = 0u;
859 rhs.m_is_closed = true;
860 rhs.m_counter = 0u;
861 }
862
863 /**
864 * Moves from rhs into this.
865 *
866 * @pre rhs is valid.
867 * @post rhs is invalid.
868 */
869 ShmControl& operator=(ShmControl&& rhs) noexcept {
871
872 m_queue = rhs.m_queue;
873 m_buffer = rhs.m_buffer;
874 m_capacity = rhs.m_capacity;
875 m_first = rhs.m_first;
876 m_last = rhs.m_last;
877 m_size = rhs.m_size;
878 m_is_closed = rhs.m_is_closed;
879 m_counter = rhs.m_counter;
880 // rhs no longer own the object.
881 rhs.m_queue = nullptr;
882 rhs.m_buffer = nullptr;
883 rhs.m_first = nullptr;
884 rhs.m_last = nullptr;
885 rhs.m_size = 0u;
886 rhs.m_is_closed = true;
887 rhs.m_counter = 0u;
888 return *this;
889 }
890
891 /**
892 * Load queue state from shared memory
893 *
894 * Performs acquire-synchronization that pairs with AtomicStoreState.
895 */
896 constexpr void AtomicLoadState() noexcept {
897 unsigned long long val = this->m_ctl_block->range.load(std::memory_order_acquire);
898 DecompressState(val);
899 }
900
901 /**
902 * Store queue range state to shared memory.
903 *
904 * Performs release-synchronization that pairs with AtomicLoadState.
905 */
906 constexpr void AtomicStoreState() noexcept {
907 unsigned long long val = CompressState();
908 this->m_ctl_block->range.store(val, std::memory_order_release);
909 }
910
911 /**
912 * Signal objects blocking on AwaitSignal.
913 */
914 void Signal() {
915 m_queue->condition.Signal();
916 }
917
918 /**
919 * Await notification signal from writer.
920 *
921 * @returns false on timeout, true otherwise.
922 */
923 template <class Rep, class Period, class Predicate>
924 [[nodiscard]] bool
925 AwaitSignal(std::chrono::duration<Rep, Period> timeout, Predicate pred) noexcept {
926 return m_queue->condition.Wait(timeout, std::move(pred));
927 }
928
929 /**
930 * Pop @a first element in queue.
931 *
932 * @note: Changes won't be seen by others until AtomicStoreState is invoked.
933 */
934 constexpr void Pop() {
936 Increment(m_first);
937 --m_size;
938 }
939
940
941 /**
942 * @returns true if queue is full, false otherwise.
943 */
944 [[nodiscard]] constexpr bool IsFull() const noexcept {
945 return m_size == this->m_capacity;
946 }
947
948 /**
949 * @returns true if queue is empty, false otherwise.
950 */
951 [[nodiscard]] constexpr bool IsEmpty() const noexcept {
952 return m_size == 0;
953 }
954
955 /**
956 * @returns true if queue is closed and will not be written to again, false otherwise.
957 * @sa Close().
958 */
959 [[nodiscard]] constexpr bool IsClosed() const noexcept {
960 return m_is_closed;
961 }
962
963 /**
964 * Close queue for further writing. Existing elements may still be written and Pop() is allowed.
965 *
966 * @note Push() will fail after this.
967 * @sa IsClosed()
968 */
969 constexpr void Close() noexcept {
970 m_is_closed = true;
971 }
972
973 /**
974 * Push element to back of queue, if `IsFull()` it will overwrite first element.
975 *
976 * @param element Element to push to queue.
977 * @returns `Error::Closed` if queue is closed.
978 * @returns `0` if element was added to queue.
979 */
980 [[nodiscard]] std::error_code Push(const T& element) noexcept {
981 /*
982 * @todo: Require that queue is not full, to enforce Popping first? Might not make a
983 * difference since were not enforcing the full protocol...
984 */
985 return Push(&element, 1);
986 }
987
988 /**
989 * @returns Error::Closed if queue is closed.
990 */
991 [[nodiscard]] std::error_code Push(const T* data, size_t n) noexcept {
992 if (m_is_closed) {
994 }
995
996 for (auto i = 0u; i < n; ++i) {
997 // Atomically store new counter value which signals that element
998 // is invalidated.
999 m_last->counter.store(m_counter++, std::memory_order_release);
1000
1001 // Copy new element to last position
1002 std::memcpy(&(m_last->element), data++, sizeof(*data));
1003
1004 if (IsFull()) {
1005 // Note: Capacity cannot be 0, IsFull() and IsEmpty() cannot both be true.
1006 //
1007 // Overwrite last element, then increment both first and last element to rotate
1008 // queue one
1009 Increment(m_last);
1010 m_first = m_last;
1011 } else {
1012 // Add without overwriting existing elements
1013 Increment(m_last);
1014 ++m_size;
1015 }
1016 }
1017 return {};
1018 }
1019
1020 [[nodiscard]] reference operator[](size_type index) {
1021 IPCQ_ASSERT(index < Size()); // check for invalid index
1022 return *Add(m_first, index);
1023 }
1024
1025 [[nodiscard]] const_reference operator[](size_type index) const {
1026 IPCQ_ASSERT(index < Size()); // check for invalid index
1027 return *Add(m_first, index);
1028 }
1029
1030 [[nodiscard]] constexpr size_type Size() const noexcept {
1031 return m_size;
1032 }
1033
1034 [[nodiscard]] constexpr value_type* GetBuffer() const noexcept {
1035 return m_buffer;
1036 }
1037
1038 [[nodiscard]] constexpr value_type* GetQueuePtr() const noexcept {
1039 return m_buffer;
1040 }
1041
1042 [[nodiscard]] constexpr ConstIterator Begin() const noexcept {
1043 return ConstIterator(this, IsEmpty() ? nullptr : m_first);
1044 }
1045
1046 [[nodiscard]] constexpr ConstIterator End() const noexcept {
1047 return ConstIterator(this, nullptr);
1048 }
1049 // @}
1050private:
1051 template <class Pointer>
1052 [[nodiscard]] constexpr Pointer Add(Pointer p, difference_type n) const noexcept {
1053 return p + (n < (m_buffer + this->m_capacity - p) ? n : n - m_capacity);
1054 }
1055
1056 template <class Pointer>
1057 constexpr void Increment(Pointer& p) const noexcept {
1058 if (++p == m_buffer + m_capacity) {
1059 p = m_buffer;
1060 }
1061 }
1062
1063 constexpr void Decrement(value_type*& p) const noexcept {
1064 if (p == m_buffer) {
1065 p = m_buffer + m_capacity;
1066 }
1067 --p;
1068 }
1069
1070 [[nodiscard]] constexpr unsigned long long CompressState() const noexcept {
1071 using namespace ipcq::detail::config;
1072 unsigned long long dest =
1073 (static_cast<unsigned long long>(m_first - m_buffer) << BIT_COUNT * 0) |
1074 (static_cast<unsigned long long>(m_last - m_buffer) << BIT_COUNT * 1) |
1075 (static_cast<unsigned long long>(m_size) << BIT_COUNT * 2) |
1076 (static_cast<unsigned long long>(m_is_closed ? BIT_IS_CLOSED : 0u));
1077 return dest;
1078 }
1079
1080 constexpr void DecompressState(unsigned long long src) noexcept {
1081 using namespace ipcq::detail::config;
1082 m_first = m_buffer + ((src >> BIT_COUNT * 0) & BIT_MASK);
1083 m_last = m_buffer + ((src >> BIT_COUNT * 1) & BIT_MASK);
1084 m_size = (src >> BIT_COUNT * 2) & BIT_MASK;
1085 m_is_closed = (src & BIT_IS_CLOSED) == BIT_IS_CLOSED;
1086 }
1087
1088 /// @name Pointers to mapped shared memory regions
1089 // @{
1090 Queue* m_queue;
1091 // @}
1092
1093 value_type* m_buffer; // Base address of queue for which offsets in m_ctl_block is valid
1094 size_t m_capacity;
1095 /**
1096 * @name Atomically updated together in the shared state with
1097 * AtomicStoreState and AtomicLoadstate.
1098 */
1099 // @{
1100 value_type* m_first; // First element
1101 value_type* m_last; // One past last element
1102 unsigned m_size; // Current size of the queue
1103 bool m_is_closed; //< Is the queue closed?
1104 // @}
1105 CounterType m_counter; // Element counter
1106};
1107
1108} // namespace ipcq::detail
1109
1110#endif
bool Wait(std::chrono::duration< Rep, Period > timeout, Predicate pred) noexcept
Definition shm.hpp:343
~BoostConditionPolicy()
Unlike std::condition_variable in C++11, it is NOT safe to invoke the destructor if all threads have ...
Definition shm.hpp:327
constexpr pid_t OwnerPid() const noexcept
Definition shm.hpp:554
ShmControlBase(ShmControlBase &&rhs) noexcept
Definition shm.hpp:503
MappedRegion m_queue_region
Mapped region to the queue (mapped as read only for Reader).
Definition shm.hpp:570
MappedRegion m_control_region
Mapped segment to the control (mapped as read only for Reader).
Definition shm.hpp:568
constexpr uint8_t FileVersion() const noexcept
Definition shm.hpp:547
char const * ShmObjectName() const
Definition shm.hpp:536
SharedMemoryObject m_shm_obj
Manages the shared memory.
Definition shm.hpp:564
ShmControlBlock * m_ctl_block
Definition shm.hpp:573
constexpr size_t Capacity() const noexcept
Definition shm.hpp:540
std::string m_topic_name
Topic name, as provided by user.
Definition shm.hpp:560
ShmControlBase(char const *name, QueueInfo queue_info, std::optional< numapp::MemPolicy > mem_policy=std::nullopt, int mmap_flags=0)
Owner constructor.
Definition shm.hpp:384
ShmControlBase & operator=(ShmControlBase &&rhs) noexcept
Definition shm.hpp:517
bool m_is_owner
Indicates if this object owns the shared memory or not.
Definition shm.hpp:576
typename ShmTraits::MappedRegion MappedRegion
Definition shm.hpp:375
typename ShmTraits::SharedMemoryObject SharedMemoryObject
Definition shm.hpp:374
ShmControlBase(char const *name)
Non-owner constructor.
Definition shm.hpp:470
ShmRemover< ShmTraits > m_shmem_remover
Definition shm.hpp:566
char const * TopicName() const noexcept
Definition shm.hpp:532
constexpr void Pop()
Pop first element in queue.
Definition shm.hpp:934
const_reference operator[](size_type index) const
Definition shm.hpp:1025
constexpr void AtomicStoreState() noexcept
Store queue range state to shared memory.
Definition shm.hpp:906
ShmControl(ShmControl &&rhs) noexcept
Moves from rhs into this.
Definition shm.hpp:842
reference operator[](size_type index)
Definition shm.hpp:1020
ShmControl(char const *name)
Open an existing (not owned) shared memory segment in read only mode.
Definition shm.hpp:777
constexpr value_type * GetBuffer() const noexcept
Definition shm.hpp:1034
typename ConstTraits< value_type >::reference const_reference
Definition shm.hpp:670
constexpr size_type Size() const noexcept
Definition shm.hpp:1030
ptrdiff_t difference_type
Definition shm.hpp:671
typename ShmControlBase< ShmTraits >::CounterType CounterType
Definition shm.hpp:660
typename NonConstTraits< value_type >::size_type size_type
Definition shm.hpp:668
detail::Iterator< ShmControl, NonConstTraits< value_type > > Iterator
Definition shm.hpp:673
constexpr value_type * GetQueuePtr() const noexcept
Definition shm.hpp:1038
constexpr bool IsEmpty() const noexcept
Definition shm.hpp:951
constexpr void AtomicLoadState() noexcept
Load queue state from shared memory.
Definition shm.hpp:896
void Signal()
Signal objects blocking on AwaitSignal.
Definition shm.hpp:914
constexpr void Close() noexcept
Close queue for further writing.
Definition shm.hpp:969
constexpr ConstIterator Begin() const noexcept
Definition shm.hpp:1042
std::error_code Push(const T &element) noexcept
Push element to back of queue, if IsFull() it will overwrite first element.
Definition shm.hpp:980
ShmControl & operator=(ShmControl &&rhs) noexcept
Moves from rhs into this.
Definition shm.hpp:869
std::error_code Push(const T *data, size_t n) noexcept
Definition shm.hpp:991
typename NonConstTraits< value_type >::reference reference
Definition shm.hpp:669
static std::error_code RemoveShmObject(char const *shm_object_name) noexcept
Attempts to remove shared memory object.
Definition shm.hpp:692
constexpr ConstIterator End() const noexcept
Definition shm.hpp:1046
bool AwaitSignal(std::chrono::duration< Rep, Period > timeout, Predicate pred) noexcept
Await notification signal from writer.
Definition shm.hpp:925
constexpr bool IsFull() const noexcept
Definition shm.hpp:944
ShmControl(char const *name, size_t capacity, std::optional< numapp::MemPolicy > mem_policy=std::nullopt, int mmap_flags=0)
Create a new (owned) named shared memory segment.
Definition shm.hpp:717
typename ShmTraits::SharedMemoryObject SharedMemoryObject
Definition shm.hpp:661
~ShmControl() noexcept
Definition shm.hpp:827
detail::Iterator< ShmControl, ConstTraits< value_type > > ConstIterator
Definition shm.hpp:672
constexpr bool IsClosed() const noexcept
Definition shm.hpp:959
Small helper that removes shared memory, to make ShmControl exception-safe with RAII.
Definition shm.hpp:237
ShmRemover & operator=(ShmRemover &&rhs) noexcept
Definition shm.hpp:254
typename ShmTraits::SharedMemoryObject SharedMemoryObject
Definition shm.hpp:239
ShmRemover(ShmRemover &&rhs) noexcept
Definition shm.hpp:248
~ShmRemover() noexcept
Definition shm.hpp:262
ShmRemover(std::string name, Policy policy)
Definition shm.hpp:243
typename ShmTraits::MappedRegion MappedRegion
Definition shm.hpp:240
Simple spinning condition.
Definition shm.hpp:278
bool Wait(std::chrono::duration< Rep, Period > timeout, Predicate pred) noexcept
Definition shm.hpp:294
Contains declarations for ipcq::Error and related.
constexpr auto BIT_MASK
Bitmask for bit count.
Definition shm.hpp:93
constexpr auto BIT_IS_CLOSED
Bit indicating if queue is closed (highest bit).
Definition shm.hpp:107
constexpr size_t TYPENAME_SEQ_LEN
Length of typename character sequences.
Definition shm.hpp:114
constexpr auto BIT_COUNT
Number of bits making up first, last and size each for a total of 63 bits.
Definition shm.hpp:87
constexpr auto MAX_VALUE
Maximum number of elements in queue.
Definition shm.hpp:100
std::error_code make_error_code(Error e)
Function overload to create std::error_code from ipcq::Error.
Definition error.hpp:120
@ None
Don't notify.
Definition writer.hpp:27
@ Closed
Queue is closed or otherwise not available for reading.
Definition error.hpp:33
@ TypeMismatch
Queue type mismatch.
Definition error.hpp:48
@ BadVersion
File version mismatch (shared object created by ipcq with different file version).
Definition error.hpp:60
@ BadFile
File type is incorrect (not an ipcq shared object).
Definition error.hpp:56
Contains declarations for shm iterator.
constexpr std::array< uint8_t, 7 > MAGIC_NUMBER
File magic number that identifies ipcq shared objects.
Definition shm.hpp:75
constexpr std::string_view SHM_OBJECT_NAME_PREFIX
Prefix used for shared object names.
Definition shm.hpp:80
void PrefaultPages(void const *begin, void const *end) noexcept
Prefault memory by reading a byte in every page starting at start up to lenght bytes.
Definition shm.hpp:164
std::string MakeShmObjectName(char const *topic_name)
Creates the shared memory object name from topic_name.
Definition shm.hpp:225
boost::posix_time::ptime AddDuration(boost::posix_time::ptime const &point, std::chrono::duration< Rep, Period > duration) noexcept
Helper that adds to posix_time::ptime a std:chrono duration.
Definition shm.hpp:213
Definition error.hpp:135
#define IPCQ_FILE_VERSION
Version number.
Definition shm.hpp:67
#define IPCQ_ASSERT(x)
Definition shm.hpp:49
boost::interprocess::shared_memory_object SharedMemoryObject
Definition shm.hpp:358
boost::interprocess::mapped_region MappedRegion
Definition shm.hpp:359
const value_type & reference
Definition traits.hpp:38
Provides standard iterator concept to ipcq queue.
Definition iterator.hpp:23
Static parameters for a given shm queue.
Definition shm.hpp:123
size_t element_offset
Offset to std::atomic<CounterType> from &Element<T>::element.
Definition shm.hpp:150
size_t element_size
Size of Element<T>.
Definition shm.hpp:154
size_t first_element_offset
Definition shm.hpp:142
size_t shm_size
Number of bytes.
Definition shm.hpp:137
static constexpr QueueInfo Make(size_t shm_size, size_t capacity, size_t first_element_offset, size_t counter_offset, size_t element_offset, size_t element_size) noexcept
Definition shm.hpp:124
size_t counter_offset
Offset to std::atomic<CounterType> from &Element<T>.
Definition shm.hpp:146
size_t capacity
Queue capacity in number of elements.
Definition shm.hpp:141
Control block residing in shared memory shared by readers and writers.
Definition shm.hpp:178
std::array< uint8_t, 7 > magic
Magic number that identifies ipcq files.
Definition shm.hpp:180
size_t element_type_hash
std::hash of element type T
Definition shm.hpp:187
std::atomic< pid_t > owner_pid
Non-zero process id of owner if owner is available or 0 if owner gracefully destructed itself but cou...
Definition shm.hpp:195
char element_type_name[config::TYPENAME_SEQ_LEN]
Definition shm.hpp:199
uint8_t version
IPCQ protocol version number (not IPCQ version number)
Definition shm.hpp:182
char condition_type_name[config::TYPENAME_SEQ_LEN]
Definition shm.hpp:200
std::atomic< unsigned long long > range
Range contains the compressed information about begin, end of queue.
Definition shm.hpp:186
size_t condition_type_hash
std::hash of ConditionPolicy
Definition shm.hpp:188
std::atomic< CounterType > counter
Definition shm.hpp:663
ConditionPolicy condition
Definition shm.hpp:680
Basic implementation detail traits.