10#ifndef IPCQ_DETAIL_SHM_HPP_
11#define IPCQ_DETAIL_SHM_HPP_
17#include <system_error>
21#include <boost/version.hpp>
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>
31#include <numapp/mempolicy.hpp>
38#define restrict __restrict__
44#define IPCQ_ENABLE_ASSERT
45#define IPCQ_ENABLE_LOG
48#ifdef IPCQ_ENABLE_ASSERT
49#define IPCQ_ASSERT(x) \
60 std::cout << x << std::endl; \
67#define IPCQ_FILE_VERSION 1
75constexpr std::array<uint8_t, 7>
MAGIC_NUMBER = {0x69, 0x70, 0x63, 0x71, 0xF0, 0x6A, 0x9C};
118namespace ipc = boost::interprocess;
157static_assert(std::is_trivially_copyable_v<QueueInfo>,
"must be trivial since it resides in SHM");
165 volatile const uint8_t* ptr =
static_cast<uint8_t const*
>(begin);
166 size_t pz = getpagesize();
167 for (; ptr < end; ptr += pz) {
186 std::atomic<unsigned long long>
range;
204static_assert(std::is_trivially_copyable<ShmControlBlock>::value,
205 "ShmControlBlock must be a POD type since it located in shared memory");
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>;
217 return point + boost::posix_time::microseconds(duration_cast<usec>(duration).count());
236template <
class ShmTraits>
245 SharedMemoryObject::remove(m_shm_name.c_str());
249 : m_shm_name(std::move(rhs.m_shm_name)), m_policy(rhs.m_policy) {
255 m_shm_name = std::move(rhs.m_shm_name);
256 m_policy = rhs.m_policy;
264 (void)SharedMemoryObject::remove(m_shm_name.c_str());
269 std::string m_shm_name;
283 m_condition.fetch_add(1, std::memory_order_release);
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();
297 auto initial = m_condition.load(std::memory_order_acquire);
299 while (initial == m_condition.load(std::memory_order_acquire)) {
300 if ((std::chrono::high_resolution_clock::now() - start) > timeout) {
303 std::this_thread::yield();
310 std::atomic<unsigned int> m_condition;
337 ipc::scoped_lock<ipc::interprocess_mutex> lock(m_mtx);
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(
347 AddDuration(boost::posix_time::microsec_clock::universal_time(), timeout),
348 std::forward<Predicate>(pred));
352 ipc::interprocess_mutex m_mtx;
353 ipc::interprocess_condition m_cond;
370template <
class ShmTraits = BoostInterprocessTraits>
386 std::optional<numapp::MemPolicy> mem_policy = std::nullopt,
395 throw std::system_error(make_error_code(Error::Overflow),
396 "Requested capacity overflows storage capacity");
417 std::optional<numapp::ScopedMemPolicy> scoped_policy;
419 scoped_policy.emplace(*mem_policy);
422#if BOOST_VERSION >= 108500
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));
446 (::boost::interprocess::default_map_options | MAP_LOCKED | MAP_POPULATE | mmap_flags)));
453 ipc::default_map_options | MAP_LOCKED | MAP_POPULATE));
473 (
std::string(
config::SHM_OBJECT_NAME_PREFIX) + name).c_str(),
484 ipc::default_map_options | MAP_LOCKED | MAP_POPULATE));
491 ipc::default_map_options | MAP_LOCKED | MAP_POPULATE));
513 rhs.m_is_owner =
false;
514 rhs.m_queue_info.shm_size = 0u;
527 rhs.m_is_owner =
false;
528 rhs.m_queue_info.shm_size = 0u;
540 [[nodiscard]]
constexpr size_t Capacity() const noexcept {
554 [[nodiscard]]
constexpr pid_t
OwnerPid() const noexcept {
657template <
class T,
class ConditionPolicy,
class ShmTraits = BoostInterprocessTraits>
695 if (!SharedMemoryObject::remove(shm_object_name)) {
696 return std::make_error_code(std::errc(errno));
719 std::optional<numapp::MemPolicy> mem_policy = std::nullopt,
726 offsetof(
Queue, buffer),
734 , m_capacity{capacity}
744 m_buffer = m_queue->buffer;
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(),
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
764 this->
m_ctl_block->owner_pid.store(getpid(), std::memory_order_release);
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;
803 if (this->
m_ctl_block->element_type_hash != std::hash<std::string>{}(
typeid(T).name())) {
808 std::hash<std::string>{}(
typeid(ConditionPolicy).name())) {
810 "ConditionPolicy type mismatch");
814 m_buffer = m_queue->
buffer;
823 "Writer exists but queue is closed");
844 , m_queue(rhs.m_queue)
845 , m_buffer(rhs.m_buffer)
846 , m_capacity(rhs.m_capacity)
847 , m_first(rhs.m_first)
850 , m_is_closed(rhs.m_is_closed)
851 , m_counter(rhs.m_counter) {
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;
859 rhs.m_is_closed =
true;
872 m_queue = rhs.m_queue;
873 m_buffer = rhs.m_buffer;
874 m_capacity = rhs.m_capacity;
875 m_first = rhs.m_first;
878 m_is_closed = rhs.m_is_closed;
879 m_counter = rhs.m_counter;
881 rhs.m_queue =
nullptr;
882 rhs.m_buffer =
nullptr;
883 rhs.m_first =
nullptr;
884 rhs.m_last =
nullptr;
886 rhs.m_is_closed =
true;
897 unsigned long long val = this->
m_ctl_block->range.load(std::memory_order_acquire);
898 DecompressState(val);
907 unsigned long long val = CompressState();
908 this->
m_ctl_block->range.store(val, std::memory_order_release);
915 m_queue->condition.Signal();
923 template <
class Rep,
class Period,
class Predicate>
925 AwaitSignal(std::chrono::duration<Rep, Period> timeout, Predicate pred)
noexcept {
926 return m_queue->condition.Wait(timeout, std::move(pred));
944 [[nodiscard]]
constexpr bool IsFull() const noexcept {
945 return m_size == this->m_capacity;
951 [[nodiscard]]
constexpr bool IsEmpty() const noexcept {
959 [[nodiscard]]
constexpr bool IsClosed() const noexcept {
980 [[nodiscard]] std::error_code
Push(
const T& element)
noexcept {
985 return Push(&element, 1);
991 [[nodiscard]] std::error_code
Push(
const T* data,
size_t n)
noexcept {
996 for (
auto i = 0u; i < n; ++i) {
999 m_last->counter.store(m_counter++, std::memory_order_release);
1002 std::memcpy(&(m_last->element), data++,
sizeof(*data));
1022 return *Add(m_first, index);
1027 return *Add(m_first, index);
1051 template <
class Po
inter>
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);
1056 template <
class Po
inter>
1057 constexpr void Increment(Pointer& p)
const noexcept {
1058 if (++p == m_buffer + m_capacity) {
1063 constexpr void Decrement(
value_type*& p)
const noexcept {
1064 if (p == m_buffer) {
1065 p = m_buffer + m_capacity;
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));
1080 constexpr void DecompressState(
unsigned long long src)
noexcept {
1081 using namespace ipcq::detail::config;
1084 m_size = (src >>
BIT_COUNT * 2) & BIT_MASK;
bool Wait(std::chrono::duration< Rep, Period > timeout, Predicate pred) noexcept
BoostConditionPolicy()=default
~BoostConditionPolicy()
Unlike std::condition_variable in C++11, it is NOT safe to invoke the destructor if all threads have ...
constexpr pid_t OwnerPid() const noexcept
ShmControlBase(ShmControlBase &&rhs) noexcept
MappedRegion m_queue_region
Mapped region to the queue (mapped as read only for Reader).
MappedRegion m_control_region
Mapped segment to the control (mapped as read only for Reader).
constexpr uint8_t FileVersion() const noexcept
char const * ShmObjectName() const
SharedMemoryObject m_shm_obj
Manages the shared memory.
ShmControlBlock * m_ctl_block
constexpr size_t Capacity() const noexcept
std::string m_topic_name
Topic name, as provided by user.
ShmControlBase(char const *name, QueueInfo queue_info, std::optional< numapp::MemPolicy > mem_policy=std::nullopt, int mmap_flags=0)
Owner constructor.
ShmControlBase & operator=(ShmControlBase &&rhs) noexcept
bool m_is_owner
Indicates if this object owns the shared memory or not.
typename ShmTraits::MappedRegion MappedRegion
typename ShmTraits::SharedMemoryObject SharedMemoryObject
ShmControlBase(char const *name)
Non-owner constructor.
ShmRemover< ShmTraits > m_shmem_remover
char const * TopicName() const noexcept
constexpr void Pop()
Pop first element in queue.
const_reference operator[](size_type index) const
constexpr void AtomicStoreState() noexcept
Store queue range state to shared memory.
ShmControl(ShmControl &&rhs) noexcept
Moves from rhs into this.
reference operator[](size_type index)
ShmControl(char const *name)
Open an existing (not owned) shared memory segment in read only mode.
constexpr value_type * GetBuffer() const noexcept
typename ConstTraits< value_type >::reference const_reference
constexpr size_type Size() const noexcept
ptrdiff_t difference_type
typename ShmControlBase< ShmTraits >::CounterType CounterType
typename NonConstTraits< value_type >::size_type size_type
detail::Iterator< ShmControl, NonConstTraits< value_type > > Iterator
constexpr value_type * GetQueuePtr() const noexcept
constexpr bool IsEmpty() const noexcept
constexpr void AtomicLoadState() noexcept
Load queue state from shared memory.
void Signal()
Signal objects blocking on AwaitSignal.
constexpr void Close() noexcept
Close queue for further writing.
constexpr ConstIterator Begin() const noexcept
std::error_code Push(const T &element) noexcept
Push element to back of queue, if IsFull() it will overwrite first element.
ShmControl & operator=(ShmControl &&rhs) noexcept
Moves from rhs into this.
std::error_code Push(const T *data, size_t n) noexcept
typename NonConstTraits< value_type >::reference reference
static std::error_code RemoveShmObject(char const *shm_object_name) noexcept
Attempts to remove shared memory object.
constexpr ConstIterator End() const noexcept
bool AwaitSignal(std::chrono::duration< Rep, Period > timeout, Predicate pred) noexcept
Await notification signal from writer.
constexpr bool IsFull() const noexcept
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.
typename ShmTraits::SharedMemoryObject SharedMemoryObject
detail::Iterator< ShmControl, ConstTraits< value_type > > ConstIterator
constexpr bool IsClosed() const noexcept
Small helper that removes shared memory, to make ShmControl exception-safe with RAII.
ShmRemover & operator=(ShmRemover &&rhs) noexcept
typename ShmTraits::SharedMemoryObject SharedMemoryObject
@ ConstructorAndDestructor
ShmRemover(ShmRemover &&rhs) noexcept
ShmRemover(std::string name, Policy policy)
typename ShmTraits::MappedRegion MappedRegion
Simple spinning condition.
bool Wait(std::chrono::duration< Rep, Period > timeout, Predicate pred) noexcept
Contains declarations for ipcq::Error and related.
constexpr auto BIT_MASK
Bitmask for bit count.
constexpr auto BIT_IS_CLOSED
Bit indicating if queue is closed (highest bit).
constexpr size_t TYPENAME_SEQ_LEN
Length of typename character sequences.
constexpr auto BIT_COUNT
Number of bits making up first, last and size each for a total of 63 bits.
constexpr auto MAX_VALUE
Maximum number of elements in queue.
std::error_code make_error_code(Error e)
Function overload to create std::error_code from ipcq::Error.
@ Closed
Queue is closed or otherwise not available for reading.
@ TypeMismatch
Queue type mismatch.
@ BadVersion
File version mismatch (shared object created by ipcq with different file version).
@ BadFile
File type is incorrect (not an ipcq shared object).
Contains declarations for shm iterator.
constexpr std::array< uint8_t, 7 > MAGIC_NUMBER
File magic number that identifies ipcq shared objects.
constexpr std::string_view SHM_OBJECT_NAME_PREFIX
Prefix used for shared object names.
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.
std::string MakeShmObjectName(char const *topic_name)
Creates the shared memory object name from topic_name.
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.
#define IPCQ_FILE_VERSION
Version number.
boost::interprocess::shared_memory_object SharedMemoryObject
boost::interprocess::mapped_region MappedRegion
const value_type & reference
Provides standard iterator concept to ipcq queue.
Static parameters for a given shm queue.
size_t element_offset
Offset to std::atomic<CounterType> from &Element<T>::element.
size_t element_size
Size of Element<T>.
size_t first_element_offset
size_t shm_size
Number of bytes.
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
size_t counter_offset
Offset to std::atomic<CounterType> from &Element<T>.
size_t capacity
Queue capacity in number of elements.
Control block residing in shared memory shared by readers and writers.
std::array< uint8_t, 7 > magic
Magic number that identifies ipcq files.
size_t element_type_hash
std::hash of element type T
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...
char element_type_name[config::TYPENAME_SEQ_LEN]
uint8_t version
IPCQ protocol version number (not IPCQ version number)
char condition_type_name[config::TYPENAME_SEQ_LEN]
std::atomic< unsigned long long > range
Range contains the compressed information about begin, end of queue.
size_t condition_type_hash
std::hash of ConditionPolicy
std::atomic< CounterType > counter
ConditionPolicy condition
Basic implementation detail traits.